|
|
|
|
@ -53,6 +53,7 @@ struct socket_node {
|
|
|
|
|
socket_t_callback_t write_cbk_sock; // For SOCK type
|
|
|
|
|
socket_callback_t except_cbk; |
|
|
|
|
void* user_data; |
|
|
|
|
const char* name; // Строковый идентификатор сокета для диагностики (литерал)
|
|
|
|
|
int active; // 1 if socket is active, 0 if freed (for reuse)
|
|
|
|
|
int enable_read; // 1 if read monitoring is enabled
|
|
|
|
|
int enable_write; // 1 if write monitoring is enabled
|
|
|
|
|
@ -73,7 +74,7 @@ struct socket_array {
|
|
|
|
|
|
|
|
|
|
static struct socket_array* socket_array_create(int initial_capacity); |
|
|
|
|
static void socket_array_destroy(struct socket_array* sa); |
|
|
|
|
static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t read_cbk, socket_callback_t write_cbk, socket_callback_t except_cbk, void* user_data); |
|
|
|
|
static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t read_cbk, socket_callback_t write_cbk, socket_callback_t except_cbk, const char* name, void* user_data); |
|
|
|
|
static int socket_array_remove(struct socket_array* sa, int fd); |
|
|
|
|
static struct socket_node* socket_array_get(struct socket_array* sa, int fd); |
|
|
|
|
|
|
|
|
|
@ -130,7 +131,7 @@ static void socket_array_destroy(struct socket_array* sa) {
|
|
|
|
|
static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t sock, int type, |
|
|
|
|
socket_callback_t read_cbk_fd, socket_callback_t write_cbk_fd, |
|
|
|
|
socket_t_callback_t read_cbk_sock, socket_t_callback_t write_cbk_sock, |
|
|
|
|
socket_callback_t except_cbk, void* user_data) { |
|
|
|
|
socket_callback_t except_cbk, const char* name, void* user_data) { |
|
|
|
|
if (!sa || fd < 0) return -1; |
|
|
|
|
// FD_SETSIZE check only for POSIX systems - Windows sockets can have any value
|
|
|
|
|
#ifndef _WIN32 |
|
|
|
|
@ -197,6 +198,7 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s
|
|
|
|
|
sa->sockets[index].write_cbk_sock = write_cbk_sock; |
|
|
|
|
sa->sockets[index].except_cbk = except_cbk; |
|
|
|
|
sa->sockets[index].user_data = user_data; |
|
|
|
|
sa->sockets[index].name = name ? name : "?"; |
|
|
|
|
sa->sockets[index].active = 1; |
|
|
|
|
sa->sockets[index].enable_read = (read_cbk_fd != NULL || read_cbk_sock != NULL) ? 1 : 0; |
|
|
|
|
sa->sockets[index].enable_write = (write_cbk_fd != NULL || write_cbk_sock != NULL) ? 1 : 0; |
|
|
|
|
@ -213,14 +215,14 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s
|
|
|
|
|
|
|
|
|
|
// Wrapper for adding regular file descriptors (pipe, file)
|
|
|
|
|
static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t read_cbk,
|
|
|
|
|
socket_callback_t write_cbk, socket_callback_t except_cbk, void* user_data) { |
|
|
|
|
socket_callback_t write_cbk, socket_callback_t except_cbk, const char* name, void* user_data) { |
|
|
|
|
return socket_array_add_internal(sa, fd, SOCKET_INVALID, SOCKET_NODE_TYPE_FD, |
|
|
|
|
read_cbk, write_cbk, NULL, NULL, except_cbk, user_data); |
|
|
|
|
read_cbk, write_cbk, NULL, NULL, except_cbk, name, user_data); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Wrapper for adding socket_t (cross-platform sockets)
|
|
|
|
|
static int socket_array_add_socket_t(struct socket_array* sa, socket_t sock, socket_t_callback_t read_cbk, |
|
|
|
|
socket_t_callback_t write_cbk, socket_callback_t except_cbk, void* user_data) { |
|
|
|
|
socket_t_callback_t write_cbk, socket_callback_t except_cbk, const char* name, void* user_data) { |
|
|
|
|
// On Windows, SOCKET is UINT_PTR, so we need to handle indexing differently
|
|
|
|
|
#ifdef _WIN32 |
|
|
|
|
int fd = (int)(intptr_t)sock; // Use socket value as index on Windows (simplified)
|
|
|
|
|
@ -230,7 +232,7 @@ static int socket_array_add_socket_t(struct socket_array* sa, socket_t sock, soc
|
|
|
|
|
if (fd < 0 || fd >= FD_SETSIZE) return -1; |
|
|
|
|
#endif |
|
|
|
|
return socket_array_add_internal(sa, fd, sock, SOCKET_NODE_TYPE_SOCK, |
|
|
|
|
NULL, NULL, read_cbk, write_cbk, except_cbk, user_data); |
|
|
|
|
NULL, NULL, read_cbk, write_cbk, except_cbk, name, user_data); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static int socket_array_remove(struct socket_array* sa, int fd) { |
|
|
|
|
@ -250,6 +252,7 @@ static int socket_array_remove(struct socket_array* sa, int fd) {
|
|
|
|
|
sa->sockets[index].write_cbk_sock = NULL; |
|
|
|
|
sa->sockets[index].except_cbk = NULL; |
|
|
|
|
sa->sockets[index].user_data = NULL; |
|
|
|
|
sa->sockets[index].name = NULL; |
|
|
|
|
sa->sockets[index].enable_read = 0; |
|
|
|
|
sa->sockets[index].enable_write = 0; |
|
|
|
|
// fd_to_index[fd] сохраняем — stale epoll события найдут неактивную ноду и пропустятся.
|
|
|
|
|
@ -297,6 +300,22 @@ static struct socket_node* socket_array_get_by_sock(struct socket_array* sa, soc
|
|
|
|
|
return &sa->sockets[index]; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Предохранитель: снять сокет с мониторинга по fd (SOCK или FD тип), если он
|
|
|
|
|
// оказался в HUP/ERR без except-обработчика. Не закрывает fd — только прекращает
|
|
|
|
|
// мониторинг, чтобы event loop не ушёл в busy-loop на вечно готовом сокете.
|
|
|
|
|
static void socket_force_unregister(struct UASYNC* ua, int fd) { |
|
|
|
|
if (!ua || fd < 0) return; |
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
if (ua->use_epoll && ua->epoll_fd >= 0) { |
|
|
|
|
epoll_ctl(ua->epoll_fd, EPOLL_CTL_DEL, fd, NULL); |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
if (socket_array_remove(ua->sockets, fd) == 0) { |
|
|
|
|
ua->socket_free_count++; |
|
|
|
|
ua->poll_fds_dirty = 1; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Callback to u_free timeout node and update counters
|
|
|
|
|
static void timeout_node_free_callback(void* user_data, void* data) { |
|
|
|
|
struct UASYNC* ua = (struct UASYNC*)user_data; |
|
|
|
|
@ -656,14 +675,14 @@ err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) {
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Instance version
|
|
|
|
|
void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, socket_callback_t write_cbk, socket_callback_t except_cbk, void* user_data) { |
|
|
|
|
void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, socket_callback_t write_cbk, socket_callback_t except_cbk, const char* name, void* user_data) { |
|
|
|
|
if (!ua || fd < 0) return NULL; |
|
|
|
|
// FD_SETSIZE check only for POSIX systems - Windows sockets can have any value
|
|
|
|
|
#ifndef _WIN32 |
|
|
|
|
if (fd >= FD_SETSIZE) return NULL; |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
int index = socket_array_add(ua->sockets, fd, read_cbk, write_cbk, except_cbk, user_data); |
|
|
|
|
int index = socket_array_add(ua->sockets, fd, read_cbk, write_cbk, except_cbk, name, user_data); |
|
|
|
|
if (index < 0) return NULL; |
|
|
|
|
|
|
|
|
|
ua->socket_alloc_count++; |
|
|
|
|
@ -688,6 +707,9 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s
|
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_SYS, "socket_add fd=%d name='%s' type=FD r=%p w=%p e=%p ud=%p", |
|
|
|
|
fd, name ? name : "?", read_cbk, write_cbk, except_cbk, user_data); |
|
|
|
|
|
|
|
|
|
// Return index-based handle
|
|
|
|
|
return (void*)(uintptr_t)(index + 1); |
|
|
|
|
} |
|
|
|
|
@ -701,6 +723,7 @@ err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) {
|
|
|
|
|
if (!node->active || node->fd < 0) return ERR_FAIL; |
|
|
|
|
|
|
|
|
|
int fd = node->fd; |
|
|
|
|
const char* rname = node->name ? node->name : "?"; |
|
|
|
|
|
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
// Remove from epoll if using epoll
|
|
|
|
|
@ -713,6 +736,7 @@ err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) {
|
|
|
|
|
if (ret == 0) { |
|
|
|
|
ua->socket_free_count++; |
|
|
|
|
ua->poll_fds_dirty = 1; // Mark poll_fds as needing rebuild
|
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_SYS, "socket_remove fd=%d name='%s'", fd, rname); |
|
|
|
|
return ERR_OK; |
|
|
|
|
} |
|
|
|
|
return ERR_FAIL; |
|
|
|
|
@ -720,11 +744,11 @@ err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) {
|
|
|
|
|
|
|
|
|
|
// Add socket_t (cross-platform socket)
|
|
|
|
|
void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t read_cbk,
|
|
|
|
|
socket_t_callback_t write_cbk, socket_t_callback_t except_cbk, void* user_data) { |
|
|
|
|
socket_t_callback_t write_cbk, socket_t_callback_t except_cbk, const char* name, void* user_data) { |
|
|
|
|
if (!ua || sock == SOCKET_INVALID) return NULL; |
|
|
|
|
|
|
|
|
|
int index = socket_array_add_socket_t(ua->sockets, sock, read_cbk, write_cbk,
|
|
|
|
|
(socket_callback_t)except_cbk, user_data); |
|
|
|
|
(socket_callback_t)except_cbk, name, user_data); |
|
|
|
|
if (index < 0) return NULL; |
|
|
|
|
|
|
|
|
|
ua->socket_alloc_count++; |
|
|
|
|
@ -751,6 +775,9 @@ void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t
|
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_SYS, "socket_add fd=%d name='%s' type=SOCK r=%p w=%p e=%p ud=%p", |
|
|
|
|
(int)sock, name ? name : "?", (void*)read_cbk, (void*)write_cbk, (void*)except_cbk, user_data); |
|
|
|
|
|
|
|
|
|
return (void*)(uintptr_t)(index + 1); /* +1: index 0 ≠ NULL */ |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -760,6 +787,7 @@ err_t uasync_remove_socket_t(struct UASYNC* ua, socket_t sock) {
|
|
|
|
|
|
|
|
|
|
struct socket_node* node = socket_array_get_by_sock(ua->sockets, sock); |
|
|
|
|
if (!node || !node->active) return ERR_FAIL; |
|
|
|
|
const char* rname = node->name ? node->name : "?"; |
|
|
|
|
|
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
if (ua->use_epoll && ua->epoll_fd >= 0) { |
|
|
|
|
@ -781,6 +809,7 @@ err_t uasync_remove_socket_t(struct UASYNC* ua, socket_t sock) {
|
|
|
|
|
if (ret == 0) { |
|
|
|
|
ua->socket_free_count++; |
|
|
|
|
ua->poll_fds_dirty = 1; |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_SYS, "socket_remove fd=%d name='%s'", fd, rname); |
|
|
|
|
return ERR_OK; |
|
|
|
|
} |
|
|
|
|
return ERR_FAIL; |
|
|
|
|
@ -991,6 +1020,17 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
|
|
|
|
|
if (events[i].events & (EPOLLERR | EPOLLHUP)) { |
|
|
|
|
if (local_except) { |
|
|
|
|
local_except(local_fd, local_ud); |
|
|
|
|
} else { |
|
|
|
|
struct socket_node* n = socket_array_get(ua->sockets, fd); |
|
|
|
|
if (!n) n = socket_array_get_by_sock(ua->sockets, fd); |
|
|
|
|
if (n && n->active && n->gen == ev_gen) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_SYS, "socket HUP/ERR without except handler: fd=%d name='%s' type=%s r=%p w=%p — unregister", |
|
|
|
|
fd, n->name ? n->name : "?", |
|
|
|
|
n->type == SOCKET_NODE_TYPE_SOCK ? "SOCK" : "FD", |
|
|
|
|
(void*)(n->type == SOCKET_NODE_TYPE_SOCK ? (void*)n->read_cbk_sock : (void*)n->read_cbk), |
|
|
|
|
(void*)(n->type == SOCKET_NODE_TYPE_SOCK ? (void*)n->write_cbk_sock : (void*)n->write_cbk)); |
|
|
|
|
socket_force_unregister(ua, fd); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
@ -1323,6 +1363,13 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
|
|
|
|
|
/* Treat as exceptional condition */ |
|
|
|
|
if (node->except_cbk) { |
|
|
|
|
node->except_cbk(node->fd, node->user_data); |
|
|
|
|
} else { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_SYS, "socket HUP/ERR/NVAL without except handler: fd=%d rev=0x%x name='%s' type=%s r=%p w=%p — unregister", |
|
|
|
|
fd, ua->poll_fds[i].revents, node->name ? node->name : "?", |
|
|
|
|
node->type == SOCKET_NODE_TYPE_SOCK ? "SOCK" : "FD", |
|
|
|
|
(void*)(node->type == SOCKET_NODE_TYPE_SOCK ? (void*)node->read_cbk_sock : (void*)node->read_cbk), |
|
|
|
|
(void*)(node->type == SOCKET_NODE_TYPE_SOCK ? (void*)node->write_cbk_sock : (void*)node->write_cbk)); |
|
|
|
|
socket_force_unregister(ua, fd); |
|
|
|
|
} |
|
|
|
|
cb_called++; |
|
|
|
|
} |
|
|
|
|
@ -1509,7 +1556,7 @@ struct UASYNC* uasync_create(void) {
|
|
|
|
|
ioctlsocket(r, FIONBIO, &mode); |
|
|
|
|
|
|
|
|
|
// Register the read socket with uasync
|
|
|
|
|
uasync_add_socket_t(ua, r, wakeup_read_callback_win, NULL, NULL, ua); // ← ua как user_data
|
|
|
|
|
uasync_add_socket_t(ua, r, wakeup_read_callback_win, NULL, NULL, "wakeup", ua); // ← ua как user_data
|
|
|
|
|
// uasync_add_socket_t(ua, r, wakeup_read_callback_win, NULL, NULL, NULL);
|
|
|
|
|
InitializeCriticalSection(&ua->posted_lock); |
|
|
|
|
|
|
|
|
|
@ -1549,7 +1596,7 @@ struct UASYNC* uasync_create(void) {
|
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
if (!ua->use_epoll) { |
|
|
|
|
uasync_add_socket(ua, ua->wakeup_pipe[0], wakeup_read_callback_posix, NULL, NULL, ua); // ← ua
|
|
|
|
|
uasync_add_socket(ua, ua->wakeup_pipe[0], wakeup_read_callback_posix, NULL, NULL, "wakeup", ua); // ← ua
|
|
|
|
|
} |
|
|
|
|
pthread_mutex_init(&ua->posted_lock, NULL); |
|
|
|
|
|
|
|
|
|
@ -1607,10 +1654,14 @@ void uasync_print_resources(struct UASYNC* ua, const char* prefix) {
|
|
|
|
|
ua->sockets->capacity, ua->sockets->count); |
|
|
|
|
for (int i = 0; i < ua->sockets->capacity; i++) { |
|
|
|
|
if (ua->sockets->sockets[i].active) { |
|
|
|
|
struct socket_node* sn = &ua->sockets->sockets[i]; |
|
|
|
|
active_sockets++; |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SYS, " Socket: fd=%d, active=%d", |
|
|
|
|
ua->sockets->sockets[i].fd, |
|
|
|
|
ua->sockets->sockets[i].active); |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SYS, " Socket: fd=%d name='%s' type=%s r=%p w=%p e=%p ud=%p", |
|
|
|
|
sn->fd, sn->name ? sn->name : "?", |
|
|
|
|
sn->type == SOCKET_NODE_TYPE_SOCK ? "SOCK" : "FD", |
|
|
|
|
(void*)(sn->type == SOCKET_NODE_TYPE_SOCK ? (void*)sn->read_cbk_sock : (void*)sn->read_cbk), |
|
|
|
|
(void*)(sn->type == SOCKET_NODE_TYPE_SOCK ? (void*)sn->write_cbk_sock : (void*)sn->write_cbk), |
|
|
|
|
(void*)sn->except_cbk, sn->user_data); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_SYS, " Total active sockets: %d", active_sockets); |
|
|
|
|
|