You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

459 lines
19 KiB

// tcp_io.c — управление TCP-соединением через uasync + ll_queue
#include "tcp_io.h"
#include "debug_config.h"
#include "mem.h"
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#ifndef _WIN32
#include <unistd.h>
#include <sys/socket.h>
#include <fcntl.h>
#endif
#ifndef MSG_NOSIGNAL
#define MSG_NOSIGNAL 0
#endif
#ifdef _WIN32
#define SHUT_WR SD_SEND
#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 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(
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)
{
if (!ua || sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: invalid args ua=%p", (void*)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->rcvbuf_size = rcvbuf_size;
tc->on_fin = on_fin;
tc->on_error = on_error;
tc->arg = arg;
tc->entry_pool = memory_pool_init(sizeof(struct ll_entry), "entry_pool");
{
size_t ds = entry_data_size > write_chunk_size ? entry_data_size : write_chunk_size;
tc->data_pool = memory_pool_init(ds, "data_pool");
}
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;
if (rcvbuf_size > 0) {
if (socket_set_buffers(sock, 0, rcvbuf_size) < 0)
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: socket_set_buffers rcvbuf=%d failed fd=%d", rcvbuf_size, (int)sock);
else
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: socket_set_buffers rcvbuf=%d ok fd=%d", rcvbuf_size, (int)sock);
}
{
struct sockaddr_storage addr;
socklen_t alen = sizeof(addr);
if (getpeername(sock, (struct sockaddr*)&addr, &alen) == 0) tc->connected = 1;
}
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: fd=%d entry=%zu chunk=%zu hw=%d lw=%d rcvbuf=%d connected=%d",
(int)sock, entry_data_size, write_chunk_size, read_high_water, read_low_water, rcvbuf_size, tc->connected);
return tc;
}
static void tcp_conn_deferred_free(void* arg) {
struct tcp_conn* tc = (struct tcp_conn*)arg;
if (tc->entry_pool) memory_pool_destroy(tc->entry_pool);
if (tc->data_pool) memory_pool_destroy(tc->data_pool);
u_free(tc);
}
void tcp_conn_destroy(struct tcp_conn* tc) {
if (!tc) return;
if (tc->destroyed) return;
tc->destroyed = 1;
DEBUG_DEBUG(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); tc->read_queue = NULL;
queue_free(tc->write_queue); tc->write_queue = NULL;
if (tc->write_buf) memory_pool_free(tc->data_pool, tc->write_buf);
if (tc->sock != SOCKET_INVALID) { socket_close_wrapper(tc->sock); tc->sock = SOCKET_INVALID; }
uasync_call_soon(tc->ua, tc, tcp_conn_deferred_free);
}
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->sock != SOCKET_INVALID) {
uasync_remove_socket_t(tc->ua, tc->sock);
tc->socket_id = NULL;
#if HAS_EPOLL
if (tc->ua && tc->ua->use_epoll && tc->ua->epoll_fd >= 0)
epoll_ctl(tc->ua->epoll_fd, EPOLL_CTL_DEL, (int)tc->sock, NULL);
#endif
socket_close_wrapper(tc->sock);
tc->sock = SOCKET_INVALID; // ДО on_error: read/write в том же epoll event увидят 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);
}
}
// ====================================================================
// Чтение из сокета
// ====================================================================
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 || !tc->read_queue) 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;
struct ll_queue* rq = tc->read_queue;
queue_data_put(rq, e);
// tc мог быть уничтожен через callback chain — read_queue обнулён
if (rq && tc->read_queue && 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_DEBUG(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;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: recv error fd=%d errno=%d", (int)tc->sock, errno);
tcp_conn_handle_error(tc, errno);
}
}
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 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) {
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 0;
}
continue;
}
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);
tcp_conn_handle_error(tc, errno);
return 1;
}
return 0;
}
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) {
if (flush_write_buf(tc)) return;
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_DEBUG(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_DEBUG(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);
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);
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);
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);
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 {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: send error fd=%d errno=%d", (int)tc->sock, errno);
memory_pool_free(tc->data_pool, e->dgram);
queue_entry_free(e);
tcp_conn_handle_error(tc, errno);
}
}
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 || !tc->write_queue) return;
if (!tc->connected) {
int err = 0;
socklen_t len = sizeof(err);
if (getsockopt(tc->sock, SOL_SOCKET, SO_ERROR, (char*)&err, &len) == 0 && err == 0) {
tc->connected = 1;
DEBUG_DEBUG(DEBUG_CATEGORY_SOCKET, "tcp_io: connect ok fd=%d", (int)tc->sock);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: connect fail fd=%d err=%d", (int)tc->sock, err);
tcp_conn_handle_error(tc, err ? err : -1);
return;
}
}
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);
if (tc->write_queue->waiter_defer) {
uasync_set_socket_write(tc->ua, tc->socket_id, 0);
tc->write_monitor = 0;
}
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;
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d connected=%d err=%d fin_r=%d fin_l=%d closed=%d wmon=%d",
(int)tc->sock, tc->connected, tc->error, tc->fin_remote, tc->fin_local, tc->closed, tc->write_monitor);
tcp_conn_handle_error(tc, -1);
}
// ====================================================================
// 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_DEBUG(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_DEBUG(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;
}