From 4d5afe5ae96608ca902980b1b2ed2706b25d8f91 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Thu, 14 May 2026 03:32:23 +0300 Subject: [PATCH] fix uasync: posted_tasks delivery on Linux, create cleanup leaks, atomic array expand --- lib/u_async.c | 144 ++++++++++++++++++++++++++++---------------------- 1 file changed, 80 insertions(+), 64 deletions(-) diff --git a/lib/u_async.c b/lib/u_async.c index 2522ca1b..d1645eef 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -129,39 +129,46 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s #ifndef _WIN32 if (fd >= FD_SETSIZE) return -1; #endif - if (fd >= sa->capacity) { - // Need to resize - double the capacity - int new_capacity = sa->capacity * 2; - if (fd >= new_capacity) new_capacity = fd + 16; // Ensure enough space - - struct socket_node* new_sockets = u_realloc(sa->sockets, new_capacity * sizeof(struct socket_node)); - int* new_fd_to_index = u_realloc(sa->fd_to_index, new_capacity * sizeof(int)); - int* new_index_to_fd = u_realloc(sa->index_to_fd, new_capacity * sizeof(int)); - int* new_active_indices = u_realloc(sa->active_indices, new_capacity * sizeof(int)); - - if (!new_sockets || !new_fd_to_index || !new_index_to_fd || !new_active_indices) { - // Allocation failed - u_free(new_sockets); - u_free(new_fd_to_index); - u_free(new_index_to_fd); - u_free(new_active_indices); - return -1; - } - - // Initialize new elements - for (int i = sa->capacity; i < new_capacity; i++) { - new_fd_to_index[i] = -1; - new_index_to_fd[i] = -1; - new_active_indices[i] = -1; - new_sockets[i].fd = -1; - new_sockets[i].active = 0; - } - - sa->sockets = new_sockets; - sa->fd_to_index = new_fd_to_index; - sa->index_to_fd = new_index_to_fd; - sa->active_indices = new_active_indices; - sa->capacity = new_capacity; + if (fd >= sa->capacity) { + int new_capacity = sa->capacity * 2; + if (fd >= new_capacity) new_capacity = fd + 16; + + struct socket_node* new_sockets = u_malloc(new_capacity * sizeof(struct socket_node)); + int* new_fd_to_index = u_malloc(new_capacity * sizeof(int)); + int* new_index_to_fd = u_malloc(new_capacity * sizeof(int)); + int* new_active_indices = u_malloc(new_capacity * sizeof(int)); + + if (!new_sockets || !new_fd_to_index || !new_index_to_fd || !new_active_indices) { + u_free(new_sockets); + u_free(new_fd_to_index); + u_free(new_index_to_fd); + u_free(new_active_indices); + return -1; + } + + memcpy(new_sockets, sa->sockets, sa->capacity * sizeof(struct socket_node)); + memcpy(new_fd_to_index, sa->fd_to_index, sa->capacity * sizeof(int)); + memcpy(new_index_to_fd, sa->index_to_fd, sa->capacity * sizeof(int)); + memcpy(new_active_indices, sa->active_indices, sa->capacity * sizeof(int)); + + for (int i = sa->capacity; i < new_capacity; i++) { + new_fd_to_index[i] = -1; + new_index_to_fd[i] = -1; + new_active_indices[i] = -1; + new_sockets[i].fd = -1; + new_sockets[i].active = 0; + } + + u_free(sa->sockets); + u_free(sa->fd_to_index); + u_free(sa->index_to_fd); + u_free(sa->active_indices); + + sa->sockets = new_sockets; + sa->fd_to_index = new_fd_to_index; + sa->index_to_fd = new_index_to_fd; + sa->active_indices = new_active_indices; + sa->capacity = new_capacity; } // Check if FD already exists @@ -839,7 +846,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events, for (int i = 0; i < n_events; i++) { if (events[i].data.fd == -1) { if (events[i].events & EPOLLIN) { - drain_wakeup_pipe(ua); + handle_wakeup(ua); } continue; } @@ -1164,7 +1171,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { /* Handle wakeup fd separately */ if (wakeup_fd_present && i == 0) { if (ua->poll_fds[i].revents & POLLIN) { - drain_wakeup_pipe(ua); + handle_wakeup(ua); } continue; } @@ -1376,27 +1383,27 @@ struct UASYNC* uasync_create(void) { #ifdef _WIN32 // Windows: self-connected UDP socket for wakeup - SOCKET r = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); - SOCKET w = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); - if (r == INVALID_SOCKET || w == INVALID_SOCKET) { - DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "Failed to create wakeup sockets"); - u_free(ua); - return NULL; - } - - struct sockaddr_in addr = {0}; - addr.sin_family = AF_INET; - addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK); - addr.sin_port = 0; - - if (bind(r, (struct sockaddr*)&addr, sizeof(addr)) == SOCKET_ERROR || - getsockname(r, (struct sockaddr*)&addr, &(int){sizeof(addr)}) == SOCKET_ERROR || - connect(w, (struct sockaddr*)&addr, sizeof(addr)) == SOCKET_ERROR) { - closesocket(r); - closesocket(w); - DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "Wakeup socket setup failed: %d", WSAGetLastError()); - u_free(ua); - return NULL; + SOCKET r = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); + SOCKET w = socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP); + if (r == INVALID_SOCKET || w == INVALID_SOCKET) { + DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "Failed to create wakeup sockets"); + if (r != INVALID_SOCKET) closesocket(r); + if (w != INVALID_SOCKET) closesocket(w); + goto create_cleanup; + } + + struct sockaddr_in addr = {0}; + addr.sin_family = AF_INET; + addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK); + addr.sin_port = 0; + + if (bind(r, (struct sockaddr*)&addr, sizeof(addr)) == SOCKET_ERROR || + getsockname(r, (struct sockaddr*)&addr, &(int){sizeof(addr)}) == SOCKET_ERROR || + connect(w, (struct sockaddr*)&addr, sizeof(addr)) == SOCKET_ERROR) { + closesocket(r); + closesocket(w); + DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "Wakeup socket setup failed: %d", WSAGetLastError()); + goto create_cleanup; } ua->wakeup_pipe[0] = (int)(intptr_t)r; @@ -1413,10 +1420,9 @@ struct UASYNC* uasync_create(void) { #else // POSIX pipe - if (pipe(ua->wakeup_pipe) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno)); - u_free(ua); - return NULL; + if (pipe(ua->wakeup_pipe) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno)); + goto create_cleanup; } ua->wakeup_initialized = 1; @@ -1428,10 +1434,20 @@ struct UASYNC* uasync_create(void) { } pthread_mutex_init(&ua->posted_lock, NULL); -#endif - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3..."); - - return ua; +#endif + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3..."); + + return ua; + +create_cleanup: + memory_pool_destroy(ua->timeout_pool); + timeout_heap_destroy(ua->timeout_heap); + socket_array_destroy(ua->sockets); +#if HAS_EPOLL + if (ua->epoll_fd >= 0) close(ua->epoll_fd); +#endif + u_free(ua); + return NULL; } // Print all resources for debugging