// tcp_io.c — управление TCP-соединением через uasync + ll_queue #include "tcp_io.h" #include "debug_config.h" #include "mem.h" #include #include #include #ifndef _WIN32 #include #include #include #endif #ifndef MSG_NOSIGNAL #define MSG_NOSIGNAL 0 #endif static void read_cb(socket_t sock, void* arg); static void write_cb(socket_t sock, void* arg); 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 const uint8_t tcp_fin_sentinel; 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, void (*on_fin)(struct tcp_conn* tc, void* arg), void (*on_error)(struct tcp_conn* tc, int err, void* arg), void* arg) { if (!ua || sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: invalid args ua=%p", ua); return NULL; } struct tcp_conn* tc = u_calloc(1, sizeof(struct tcp_conn)); if (!tc) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: u_calloc failed"); return NULL; } tc->ua = ua; tc->sock = sock; tc->entry_data_size = entry_data_size; tc->write_chunk_size = write_chunk_size; tc->read_high_water = read_high_water; tc->read_low_water = read_low_water; tc->on_fin = on_fin; tc->on_error = on_error; tc->arg = arg; tc->entry_pool = memory_pool_init(sizeof(struct ll_entry)); { size_t ds = entry_data_size > write_chunk_size ? entry_data_size : write_chunk_size; tc->data_pool = memory_pool_init(ds); } if (!tc->entry_pool || !tc->data_pool) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: memory_pool_init failed"); if (tc->entry_pool) memory_pool_destroy(tc->entry_pool); if (tc->data_pool) memory_pool_destroy(tc->data_pool); u_free(tc); return NULL; } tc->read_queue = queue_new(ua, 0, 0, 0, "tcp_io_rq"); tc->write_queue = queue_new(ua, 0, 0, 0, "tcp_io_wq"); if (!tc->read_queue || !tc->write_queue) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: queue_new failed"); if (tc->read_queue) queue_free(tc->read_queue); if (tc->write_queue) queue_free(tc->write_queue); memory_pool_destroy(tc->entry_pool); memory_pool_destroy(tc->data_pool); u_free(tc); return NULL; } queue_set_threshold(tc->read_queue, read_low_water, 0); queue_set_threshold(tc->write_queue, 32, 0); queue_set_callback(tc->write_queue, write_queue_fetch_cb, tc); queue_set_waiter_defer(tc->write_queue, 1); tc->socket_id = uasync_add_socket_t(ua, sock, read_cb, write_cb, error_cb, tc); if (!tc->socket_id) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: uasync_add_socket_t failed"); queue_free(tc->read_queue); queue_free(tc->write_queue); memory_pool_destroy(tc->entry_pool); memory_pool_destroy(tc->data_pool); u_free(tc); return NULL; } tc->write_monitor = 1; { struct sockaddr_storage addr; socklen_t alen = sizeof(addr); if (getpeername(sock, (struct sockaddr*)&addr, &alen) == 0) tc->connected = 1; } DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: fd=%d entry=%zu chunk=%zu hw=%d lw=%d connected=%d", (int)sock, entry_data_size, write_chunk_size, read_high_water, read_low_water, tc->connected); return tc; } void tcp_conn_destroy(struct tcp_conn* tc) { if (!tc) return; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_conn_destroy: fd=%d connected=%d error=%d fin_remote=%d fin_local=%d closed=%d", (int)tc->sock, tc->connected, tc->error, tc->fin_remote, tc->fin_local, tc->closed); if (tc->socket_id) { uasync_remove_socket_t(tc->ua, tc->sock); tc->socket_id = NULL; } 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); queue_set_callback(tc->write_queue, NULL, NULL); struct ll_entry* e; while ((e = queue_data_get(tc->read_queue)) != NULL) { if (e->dgram) memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(tc->read_queue); } while ((e = queue_data_get(tc->write_queue)) != NULL) { if (e->dgram && e->dgram != &tcp_fin_sentinel) memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(tc->write_queue); } queue_free(tc->read_queue); queue_free(tc->write_queue); if (tc->write_buf) memory_pool_free(tc->data_pool, tc->write_buf); if (tc->sock != SOCKET_INVALID) socket_close_wrapper(tc->sock); memory_pool_destroy(tc->entry_pool); memory_pool_destroy(tc->data_pool); u_free(tc); } // ==================================================================== // Чтение из сокета // ==================================================================== static void read_cb(socket_t sock, void* arg) { (void)sock; struct tcp_conn* tc = (struct tcp_conn*)arg; if (!tc || tc->sock == SOCKET_INVALID) return; struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool); uint8_t* buf = memory_pool_alloc(tc->data_pool); if (!e || !buf) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: read_cb alloc failed fd=%d", (int)tc->sock); if (e) queue_entry_free(e); if (buf) memory_pool_free(tc->data_pool, buf); uasync_set_socket_read(tc->ua, tc->socket_id, 0); tc->read_paused = 1; return; } 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); if (tc->read_queue->count >= tc->read_high_water && !tc->read_paused) { uasync_set_socket_read(tc->ua, tc->socket_id, 0); tc->read_paused = 1; queue_waiter_wait(tc->read_queue, &tc->read_waiter, resume_read_cb, tc); } } else if (n == 0) { memory_pool_free(tc->data_pool, buf); queue_entry_free(e); tc->fin_remote = 1; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: FIN fd=%d", (int)tc->sock); uasync_set_socket_read(tc->ua, tc->socket_id, 0); queue_waiter_cancel(tc->read_queue, &tc->read_waiter); if (tc->read_queue->count == 0) { if (tc->on_fin) tc->on_fin(tc, tc->arg); } else { queue_set_empty_callback(tc->read_queue, fin_deferred_cb, tc); } } else { 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); } } static void resume_read_cb(struct ll_queue* q, void* arg) { (void)q; struct tcp_conn* tc = (struct tcp_conn*)arg; if (!tc || tc->sock == SOCKET_INVALID) return; uasync_set_socket_read(tc->ua, tc->socket_id, 1); tc->read_paused = 0; } static void fin_deferred_cb(struct ll_queue* q, void* arg) { (void)q; struct tcp_conn* tc = (struct tcp_conn*)arg; if (tc->on_fin) tc->on_fin(tc, tc->arg); } // ==================================================================== // Запись в сокет // ==================================================================== static void 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) { tc->write_offset += (size_t)n; if (tc->write_offset >= tc->write_len) { memory_pool_free(tc->data_pool, tc->write_buf); tc->write_buf = NULL; tc->write_len = 0; tc->write_offset = 0; DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: flush_wbuf done fd=%d", (int)tc->sock); return; } continue; } if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) return; 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); return; } } static void write_queue_fetch_cb(struct ll_queue* q, void* arg) { struct tcp_conn* tc = (struct tcp_conn*)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 (tc->write_buf) return; // остался остаток, ждём write_cb } if (!tc->connected) { DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch skip fd=%d not connected", (int)tc->sock); return; } struct ll_entry* e = queue_data_get(q); if (!e) { uasync_set_socket_write(tc->ua, tc->socket_id, 0); tc->write_monitor = 0; DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT OFF fd=%d (queue empty)", (int)tc->sock); queue_resume_callback(q); if (tc->on_flushed) { void (*cb)(struct tcp_conn*, void*) = tc->on_flushed; tc->on_flushed = NULL; cb(tc, tc->arg); } return; } if (e->len == 0) { if (!e->dgram) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: CLOSE fd=%d", (int)tc->sock); queue_entry_free(e); queue_resume_callback(q); uasync_remove_socket_t(tc->ua, tc->sock); socket_close_wrapper(tc->sock); tc->socket_id = NULL; tc->sock = SOCKET_INVALID; tc->closed = 1; if (tc->on_closed) tc->on_closed(tc, tc->arg); } else if (e->dgram == &tcp_fin_sentinel) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: FIN fd=%d", (int)tc->sock); queue_entry_free(e); queue_resume_callback(q); shutdown(tc->sock, SHUT_WR); tc->fin_local = 1; if (tc->on_fin_sent) tc->on_fin_sent(tc, tc->arg); } else { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch zero-len entry with unknown dgram=%p fd=%d", (void*)e->dgram, (int)tc->sock); queue_entry_free(e); queue_resume_callback(q); } return; } ssize_t n = send(tc->sock, e->dgram, e->len, MSG_NOSIGNAL); if (n == (ssize_t)e->len) { memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); queue_resume_callback(q); } else if (n > 0) { size_t rem = e->len - (size_t)n; 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); return; } memcpy(tc->write_buf, e->dgram + n, rem); tc->write_len = rem; tc->write_offset = 0; memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); uasync_set_socket_write(tc->ua, tc->socket_id, 1); tc->write_monitor = 1; DEBUG_TRACE(DEBUG_CATEGORY_SOCKET, "tcp_io: wq_fetch POLLOUT ON fd=%d (partial send, wbuf=%zu bytes)", (int)tc->sock, tc->write_len); } else if (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)) { 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); return; } memcpy(tc->write_buf, e->dgram, e->len); tc->write_len = e->len; tc->write_offset = 0; memory_pool_free(tc->data_pool, e->dgram); queue_entry_free(e); uasync_set_socket_write(tc->ua, tc->socket_id, 1); 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); } } static void write_cb(socket_t sock, void* arg) { (void)sock; struct tcp_conn* tc = (struct tcp_conn*)arg; if (!tc || tc->sock == SOCKET_INVALID) return; if (!tc->connected) { int err = 0; socklen_t len = sizeof(err); if (getsockopt(tc->sock, SOL_SOCKET, SO_ERROR, &err, &len) == 0 && err == 0) { 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); return; } } flush_write_buf(tc); 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); } } // ==================================================================== // Асинхронная ошибка сокета // ==================================================================== 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); } // ==================================================================== // FIN / Close через очередь отправки // ==================================================================== int tcp_conn_push_fin(struct tcp_conn* tc) { if (!tc || tc->sock == SOCKET_INVALID) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin invalid tc"); return -1; } if (tc->fin_local) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin already fin_local fd=%d", (int)tc->sock); return -1; } if (tc->closed) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin already closed fd=%d", (int)tc->sock); return -1; } struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin entry_pool exhausted fd=%d", (int)tc->sock); return -1; } e->dgram = (uint8_t*)&tcp_fin_sentinel; e->len = 0; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: push_fin fd=%d wq=%d", (int)tc->sock, tc->write_queue->count + 1); queue_data_put(tc->write_queue, e); return 0; } int tcp_conn_push_close(struct tcp_conn* tc) { if (!tc || tc->sock == SOCKET_INVALID) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close invalid tc"); return -1; } if (tc->closed) { DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close already closed fd=%d", (int)tc->sock); return -1; } struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close entry_pool exhausted fd=%d", (int)tc->sock); return -1; } e->dgram = NULL; e->len = 0; DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: push_close fd=%d wq=%d", (int)tc->sock, tc->write_queue->count + 1); queue_data_put(tc->write_queue, e); return 0; } void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_conn* tc, void* arg)) { if (!tc) return; tc->on_flushed = on_flushed; } void tcp_conn_pause_read(struct tcp_conn* tc) { if (!tc || tc->sock == SOCKET_INVALID) return; uasync_set_socket_read(tc->ua, tc->socket_id, 0); queue_waiter_cancel(tc->read_queue, &tc->read_waiter); tc->read_paused = 1; }