diff --git a/lib/u_async.c b/lib/u_async.c index b86623e8..042c08ab 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -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); diff --git a/lib/u_async.h b/lib/u_async.h index de289104..f78ac493 100644 --- a/lib/u_async.h +++ b/lib/u_async.h @@ -113,9 +113,9 @@ void* uasync_call_soon(struct UASYNC* ua, void* user_arg, timeout_callback_t cal err_t uasync_call_soon_cancel(struct UASYNC* ua, void* t_id); // Sockets - for regular file descriptors (pipe, file) -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_arg); +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_arg); // Sockets - for socket_t (cross-platform sockets) -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_arg); +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, const char* name, void* user_arg); err_t uasync_remove_socket(struct UASYNC* ua, void* s_id); err_t uasync_remove_socket_t(struct UASYNC* ua, socket_t sock);