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.
331 lines
12 KiB
331 lines
12 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 |
|
|
|
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); |
|
|
|
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=%d", |
|
(int)tc->sock, tc->connected, tc->error, tc->fin); |
|
|
|
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) 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); |
|
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 = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_io: FIN fd=%d", (int)tc->sock); |
|
uasync_set_socket_read(tc->ua, tc->socket_id, 0); |
|
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; |
|
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->fin) return; |
|
|
|
if (tc->write_buf) { |
|
flush_write_buf(tc); |
|
if (tc->write_buf) return; // остался остаток, ждём write_cb |
|
} |
|
|
|
if (!tc->connected) return; // ждём connect |
|
|
|
struct ll_entry* e = queue_data_get(q); |
|
if (!e) { |
|
uasync_set_socket_write(tc->ua, tc->socket_id, 0); |
|
tc->write_monitor = 0; |
|
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; |
|
} |
|
|
|
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; |
|
} 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; |
|
} 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_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) 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_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: async error fd=%d", (int)tc->sock); |
|
if (tc->on_error) tc->on_error(tc, -1, tc->arg); |
|
} |
|
|
|
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; |
|
}
|
|
|