From a97a38d867d273d55a8251b5cfd7ef2cceff2937 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Thu, 2 Jul 2026 00:43:37 +0300 Subject: [PATCH] fix: set tc->sock=SOCKET_INVALID before on_error in handle_error MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit В одном epoll_wait batch три события для fd: EPOLLERR (освобождает tc), EPOLLIN (read_cb), EPOLLOUT (write_cb). read_cb/write_cb guard проверяет tc->sock==SOCKET_INVALID, но sock обнулялся ПОСЛЕ on_error. Перенос sock=SOCKET_INVALID до on_error защищает от use-after-free. Доказано: без фикса segfault на tcp_io.c:374 (write_cb → tc->connected). С фиксом тест работает без крашей. --- lib/tcp_io.c | 3 +- lib/u_async.c | 1242 ++++++++++++++++--------------- tests/test_uasync_socket_race.c | 8 +- 3 files changed, 629 insertions(+), 624 deletions(-) diff --git a/lib/tcp_io.c b/lib/tcp_io.c index effe2651..0080386e 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -158,7 +158,6 @@ static void tcp_conn_handle_error(struct tcp_conn* tc, int err) } if (tc->sock != SOCKET_INVALID) { - // Гарантированно удаляем из epoll (uasync_remove_socket_t может пропустить DEL если нода уже inactive) uasync_remove_socket_t(tc->ua, tc->sock); tc->socket_id = NULL; #if HAS_EPOLL @@ -166,7 +165,7 @@ static void tcp_conn_handle_error(struct tcp_conn* tc, int err) epoll_ctl(tc->ua->epoll_fd, EPOLL_CTL_DEL, (int)tc->sock, NULL); #endif socket_close_wrapper(tc->sock); - tc->sock = SOCKET_INVALID; + tc->sock = SOCKET_INVALID; // ДО on_error: read/write в том же epoll event увидят INVALID } queue_set_callback(tc->write_queue, NULL, NULL); diff --git a/lib/u_async.c b/lib/u_async.c index 2128986f..2ee24afb 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -1,9 +1,9 @@ // uasync.c -#include "u_async.h" -#include "platform_compat.h" -#include "debug_config.h" -#include "mem.h" +#include "u_async.h" +#include "platform_compat.h" +#include "debug_config.h" +#include "mem.h" #include "memory_pool.h" #include #include @@ -30,44 +30,44 @@ -// Timeout node with safe cancellation -struct timeout_node { - char name[16]; - void* arg; - timeout_callback_t callback; - uint64_t expiration_ms; // absolute expiration time in milliseconds - struct UASYNC* ua; // Pointer back to uasync instance for counter updates - struct timeout_node* next; // For immediate queue (FIFO) - size_t heap_index; // Position in timeout_heap (SIZE_MAX if not in heap) +// Timeout node with safe cancellation +struct timeout_node { + char name[16]; + void* arg; + timeout_callback_t callback; + uint64_t expiration_ms; // absolute expiration time in milliseconds + struct UASYNC* ua; // Pointer back to uasync instance for counter updates + struct timeout_node* next; // For immediate queue (FIFO) + size_t heap_index; // Position in timeout_heap (SIZE_MAX if not in heap) }; // Socket node with array-based storage -struct socket_node { - int fd; // File descriptor (for pipe, file) - socket_t sock; // Socket (for cross-platform sockets) - int type; // SOCKET_NODE_TYPE_FD or SOCKET_NODE_TYPE_SOCK - socket_callback_t read_cbk; // For FD type - socket_callback_t write_cbk; // For FD type - socket_t_callback_t read_cbk_sock; // For SOCK type - socket_t_callback_t write_cbk_sock; // For SOCK type - socket_callback_t except_cbk; - void* user_data; - 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 - uint16_t gen; // generation counter for epoll event validation +struct socket_node { + int fd; // File descriptor (for pipe, file) + socket_t sock; // Socket (for cross-platform sockets) + int type; // SOCKET_NODE_TYPE_FD or SOCKET_NODE_TYPE_SOCK + socket_callback_t read_cbk; // For FD type + socket_callback_t write_cbk; // For FD type + socket_t_callback_t read_cbk_sock; // For SOCK type + socket_t_callback_t write_cbk_sock; // For SOCK type + socket_callback_t except_cbk; + void* user_data; + 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 + uint16_t gen; // generation counter for epoll event validation }; // Array-based socket management for O(1) operations -struct socket_array { - struct socket_node* sockets; // Dynamic array of socket nodes - int* fd_to_index; // FD to array index mapping - int* index_to_fd; // Array index to FD mapping - int* active_indices; // Array of indices of active sockets (for O(1) traversal) - int capacity; // Total allocated capacity - int count; // Number of active sockets - int max_fd; // Maximum FD for bounds checking - uint16_t gen_counter; // incrementing generation for epoll stale-event detection +struct socket_array { + struct socket_node* sockets; // Dynamic array of socket nodes + int* fd_to_index; // FD to array index mapping + int* index_to_fd; // Array index to FD mapping + int* active_indices; // Array of indices of active sockets (for O(1) traversal) + int capacity; // Total allocated capacity + int count; // Number of active sockets + int max_fd; // Maximum FD for bounds checking + uint16_t gen_counter; // incrementing generation for epoll stale-event detection }; static struct socket_array* socket_array_create(int initial_capacity); @@ -82,8 +82,8 @@ static struct socket_node* socket_array_get(struct socket_array* sa, int fd); static struct socket_array* socket_array_create(int initial_capacity) { if (initial_capacity < 4) initial_capacity = 4; // Minimum capacity - struct socket_array* sa = u_malloc(sizeof(struct socket_array)); - if (!sa) return NULL; + struct socket_array* sa = u_malloc(sizeof(struct socket_array)); + if (!sa) return NULL; sa->gen_counter = 0; sa->sockets = u_calloc(initial_capacity, sizeof(struct socket_node)); @@ -170,20 +170,20 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s sa->capacity = new_capacity; } - // Check if FD already has a node — reuse inactive slot if present - int index; - int existing = sa->fd_to_index[fd]; - if (existing != -1) { - if (sa->sockets[existing].active) return -1; // FD уже занят активной нодой — ошибка - index = existing; // переиспользуем неактивный слот (после socket_array_remove) - } else { - // Find first free slot - index = -1; - for (int i = 0; i < sa->capacity; i++) { - if (!sa->sockets[i].active) { index = i; break; } - } - if (index == -1) return -1; // No free slots - sa->fd_to_index[fd] = index; // новая привязка fd→слот + // Check if FD already has a node — reuse inactive slot if present + int index; + int existing = sa->fd_to_index[fd]; + if (existing != -1) { + if (sa->sockets[existing].active) return -1; // FD уже занят активной нодой — ошибка + index = existing; // переиспользуем неактивный слот (после socket_array_remove) + } else { + // Find first free slot + index = -1; + for (int i = 0; i < sa->capacity; i++) { + if (!sa->sockets[i].active) { index = i; break; } + } + if (index == -1) return -1; // No free slots + sa->fd_to_index[fd] = index; // новая привязка fd→слот } // Add the socket @@ -194,13 +194,13 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s sa->sockets[index].write_cbk = write_cbk_fd; sa->sockets[index].read_cbk_sock = read_cbk_sock; 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].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; - sa->sockets[index].gen = ++sa->gen_counter; - + sa->sockets[index].except_cbk = except_cbk; + sa->sockets[index].user_data = user_data; + 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; + sa->sockets[index].gen = ++sa->gen_counter; + sa->index_to_fd[index] = fd; sa->active_indices[sa->count] = index; // Add to active list sa->count++; @@ -238,21 +238,21 @@ static int socket_array_remove(struct socket_array* sa, int fd) { int index = sa->fd_to_index[fd]; if (index == -1 || !sa->sockets[index].active) return -1; // FD not found - // Mark as inactive and clear all pointers - sa->sockets[index].active = 0; - sa->sockets[index].fd = -1; - sa->sockets[index].sock = SOCKET_INVALID; - sa->sockets[index].type = SOCKET_NODE_TYPE_FD; - sa->sockets[index].read_cbk = NULL; - sa->sockets[index].write_cbk = NULL; - sa->sockets[index].read_cbk_sock = NULL; - sa->sockets[index].write_cbk_sock = NULL; - sa->sockets[index].except_cbk = NULL; - sa->sockets[index].user_data = NULL; - sa->sockets[index].enable_read = 0; - sa->sockets[index].enable_write = 0; - // fd_to_index[fd] сохраняем — stale epoll события найдут неактивную ноду и пропустятся. - // При переиспользовании fd, socket_array_add_internal перезапишет этот же слот. + // Mark as inactive and clear all pointers + sa->sockets[index].active = 0; + sa->sockets[index].fd = -1; + sa->sockets[index].sock = SOCKET_INVALID; + sa->sockets[index].type = SOCKET_NODE_TYPE_FD; + sa->sockets[index].read_cbk = NULL; + sa->sockets[index].write_cbk = NULL; + sa->sockets[index].read_cbk_sock = NULL; + sa->sockets[index].write_cbk_sock = NULL; + sa->sockets[index].except_cbk = NULL; + sa->sockets[index].user_data = NULL; + sa->sockets[index].enable_read = 0; + sa->sockets[index].enable_write = 0; + // fd_to_index[fd] сохраняем — stale epoll события найдут неактивную ноду и пропустятся. + // При переиспользовании fd, socket_array_add_internal перезапишет этот же слот. sa->index_to_fd[index] = -1; // Remove from active_indices by swapping with last element @@ -296,13 +296,13 @@ static struct socket_node* socket_array_get_by_sock(struct socket_array* sa, soc return &sa->sockets[index]; } -// 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; - struct timeout_node* node = (struct timeout_node*)data; - (void)node; // Not used directly, but keep for consistency - ua->timer_free_count++; - memory_pool_free(ua->timeout_pool, data); +// 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; + struct timeout_node* node = (struct timeout_node*)data; + (void)node; // Not used directly, but keep for consistency + ua->timer_free_count++; + memory_pool_free(ua->timeout_pool, data); } // Helper to get current time @@ -339,27 +339,27 @@ uint64_t get_time_tb(void) { // QueryPerformanceCounter(&count); // Получаем текущее значение счётчика // return (uint64_t)(count.QuadPart * 10000ULL) / (uint64_t)freq.QuadPart; // Преобразуем в требуемые единицы времени //} -#else -uint64_t get_time_tb(void) { - struct timespec ts; - clock_gettime(CLOCK_MONOTONIC, &ts); - return (uint64_t)ts.tv_sec * 10000ULL + (uint64_t)ts.tv_nsec / 100000ULL; // Преобразуем в требуемые единицы времени -} -#endif - -#ifdef _WIN32 -uint64_t get_time_us(void) { - LARGE_INTEGER freq, count; - QueryPerformanceFrequency(&freq); - QueryPerformanceCounter(&count); - return (uint64_t)(count.QuadPart * 1000000ULL / freq.QuadPart); -} -#else -uint64_t get_time_us(void) { - struct timespec ts; - clock_gettime(CLOCK_MONOTONIC, &ts); - return (uint64_t)ts.tv_sec * 1000000ULL + (uint64_t)ts.tv_nsec / 1000ULL; -} +#else +uint64_t get_time_tb(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (uint64_t)ts.tv_sec * 10000ULL + (uint64_t)ts.tv_nsec / 100000ULL; // Преобразуем в требуемые единицы времени +} +#endif + +#ifdef _WIN32 +uint64_t get_time_us(void) { + LARGE_INTEGER freq, count; + QueryPerformanceFrequency(&freq); + QueryPerformanceCounter(&count); + return (uint64_t)(count.QuadPart * 1000000ULL / freq.QuadPart); +} +#else +uint64_t get_time_us(void) { + struct timespec ts; + clock_gettime(CLOCK_MONOTONIC, &ts); + return (uint64_t)ts.tv_sec * 1000000ULL + (uint64_t)ts.tv_nsec / 1000ULL; +} #endif @@ -443,56 +443,56 @@ static uint64_t timeval_to_ms(const struct timeval* tv) { // Simplified timeout handling without reference counting -// Process expired timeouts with safe cancellation -static void process_timeouts(struct UASYNC* ua) { - if (!ua) return; - - // Сначала обрабатываем immediate_queue (FIFO) - while (ua->immediate_queue_head) { - struct timeout_node* node = ua->immediate_queue_head; - ua->immediate_queue_head = node->next; - if (!ua->immediate_queue_head) { - ua->immediate_queue_tail = NULL; - } - - if (node && node->callback) { - DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "timer→immediate %s", node->name[0] ? node->name : ""); - node->callback(node->arg); - } - - if (node && node->ua) { - node->ua->timer_free_count++; - } - memory_pool_free(ua->timeout_pool, node); - } - - if (!ua->timeout_heap) return; - - struct timeval now_tv; - get_current_time(&now_tv); - uint64_t now_ms = timeval_to_ms(&now_tv); - - while (1) { - TimeoutEntry entry; - if (timeout_heap_peek(ua->timeout_heap, &entry) != 0) break; - if (entry.expiration > now_ms) break; - - // Pop the expired timeout - timeout_heap_pop(ua->timeout_heap, &entry); - struct timeout_node* node = (struct timeout_node*)entry.data; - - if (node && node->callback) { - DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "timer→%s expired", node->name[0] ? node->name : ""); - node->callback(node->arg); - } - - // Always u_free the node after processing - if (node && node->ua) { - node->ua->timer_free_count++; - } - memory_pool_free(ua->timeout_pool, node); - continue; // Process next expired timeout - } +// Process expired timeouts with safe cancellation +static void process_timeouts(struct UASYNC* ua) { + if (!ua) return; + + // Сначала обрабатываем immediate_queue (FIFO) + while (ua->immediate_queue_head) { + struct timeout_node* node = ua->immediate_queue_head; + ua->immediate_queue_head = node->next; + if (!ua->immediate_queue_head) { + ua->immediate_queue_tail = NULL; + } + + if (node && node->callback) { + DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "timer→immediate %s", node->name[0] ? node->name : ""); + node->callback(node->arg); + } + + if (node && node->ua) { + node->ua->timer_free_count++; + } + memory_pool_free(ua->timeout_pool, node); + } + + if (!ua->timeout_heap) return; + + struct timeval now_tv; + get_current_time(&now_tv); + uint64_t now_ms = timeval_to_ms(&now_tv); + + while (1) { + TimeoutEntry entry; + if (timeout_heap_peek(ua->timeout_heap, &entry) != 0) break; + if (entry.expiration > now_ms) break; + + // Pop the expired timeout + timeout_heap_pop(ua->timeout_heap, &entry); + struct timeout_node* node = (struct timeout_node*)entry.data; + + if (node && node->callback) { + DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "timer→%s expired", node->name[0] ? node->name : ""); + node->callback(node->arg); + } + + // Always u_free the node after processing + if (node && node->ua) { + node->ua->timer_free_count++; + } + memory_pool_free(ua->timeout_pool, node); + continue; // Process next expired timeout + } } // Compute time to next timeout @@ -527,118 +527,118 @@ static void get_next_timeout(struct UASYNC* ua, struct timeval* tv) { -// Instance version -void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_callback_t callback, const char* name) { - if (!ua || timeout_tb < 0 || !callback) return NULL; - if (!ua->timeout_heap) return NULL; - -// DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: timeout=%d.%d ms, arg=%p, callback=%p", timeout_tb/10, timeout_tb%10, arg, callback); - - struct timeout_node* node = memory_pool_alloc(ua->timeout_pool); - if (!node) { - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: failed to allocate node"); - return NULL; - } - ua->timer_alloc_count++; - - if (name) { - strncpy(node->name, name, sizeof(node->name) - 1); - node->name[sizeof(node->name) - 1] = '\0'; - } else { - node->name[0] = '\0'; - } - node->arg = arg; - node->callback = callback; - node->ua = ua; - node->heap_index = SIZE_MAX; - - // Calculate expiration time in milliseconds - struct timeval now; - get_current_time(&now); - timeval_add_tb(&now, timeout_tb); - node->expiration_ms = timeval_to_ms(&now); - - // Add to heap +// Instance version +void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_callback_t callback, const char* name) { + if (!ua || timeout_tb < 0 || !callback) return NULL; + if (!ua->timeout_heap) return NULL; + +// DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: timeout=%d.%d ms, arg=%p, callback=%p", timeout_tb/10, timeout_tb%10, arg, callback); + + struct timeout_node* node = memory_pool_alloc(ua->timeout_pool); + if (!node) { + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: failed to allocate node"); + return NULL; + } + ua->timer_alloc_count++; + + if (name) { + strncpy(node->name, name, sizeof(node->name) - 1); + node->name[sizeof(node->name) - 1] = '\0'; + } else { + node->name[0] = '\0'; + } + node->arg = arg; + node->callback = callback; + node->ua = ua; + node->heap_index = SIZE_MAX; + + // Calculate expiration time in milliseconds + struct timeval now; + get_current_time(&now); + timeval_add_tb(&now, timeout_tb); + node->expiration_ms = timeval_to_ms(&now); + + // Add to heap if (timeout_heap_push(ua->timeout_heap, node->expiration_ms, node, &node->heap_index) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: failed to push to heap"); - memory_pool_free(ua->timeout_pool, node); - ua->timer_free_count++; // Balance the alloc counter - return NULL; - } - - return node; -} - -// Immediate execution in next mainloop (FIFO order) -void* uasync_call_soon(struct UASYNC* ua, void* user_arg, timeout_callback_t callback) { - if (!ua || !callback) return NULL; - if (!ua->timeout_pool) return NULL; - - struct timeout_node* node = memory_pool_alloc(ua->timeout_pool); - if (!node) { - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_call_soon: failed to allocate node"); - return NULL; - } - ua->timer_alloc_count++; - - node->name[0] = '\0'; - node->arg = user_arg; - node->callback = callback; - node->ua = ua; - node->expiration_ms = 0; - node->next = NULL; - node->heap_index = SIZE_MAX; - - // FIFO: добавляем в конец очереди - if (ua->immediate_queue_tail) { - ua->immediate_queue_tail->next = node; - ua->immediate_queue_tail = node; - } else { - ua->immediate_queue_head = ua->immediate_queue_tail = node; - } - - return node; -} - -// Cancel immediate callback by setting callback to NULL - O(1) -err_t uasync_call_soon_cancel(struct UASYNC* ua, void* t_id) { - if (!ua || !t_id) return ERR_FAIL; - - struct timeout_node* node = (struct timeout_node*)t_id; - if (node->ua != ua) return ERR_FAIL; - - // Simply nullify callback - will be skipped in process_timeouts - node->callback = NULL; - - return ERR_OK; -} - - - + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: failed to push to heap"); + memory_pool_free(ua->timeout_pool, node); + ua->timer_free_count++; // Balance the alloc counter + return NULL; + } + + return node; +} + +// Immediate execution in next mainloop (FIFO order) +void* uasync_call_soon(struct UASYNC* ua, void* user_arg, timeout_callback_t callback) { + if (!ua || !callback) return NULL; + if (!ua->timeout_pool) return NULL; + + struct timeout_node* node = memory_pool_alloc(ua->timeout_pool); + if (!node) { + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_call_soon: failed to allocate node"); + return NULL; + } + ua->timer_alloc_count++; + + node->name[0] = '\0'; + node->arg = user_arg; + node->callback = callback; + node->ua = ua; + node->expiration_ms = 0; + node->next = NULL; + node->heap_index = SIZE_MAX; + + // FIFO: добавляем в конец очереди + if (ua->immediate_queue_tail) { + ua->immediate_queue_tail->next = node; + ua->immediate_queue_tail = node; + } else { + ua->immediate_queue_head = ua->immediate_queue_tail = node; + } + + return node; +} + +// Cancel immediate callback by setting callback to NULL - O(1) +err_t uasync_call_soon_cancel(struct UASYNC* ua, void* t_id) { + if (!ua || !t_id) return ERR_FAIL; + + struct timeout_node* node = (struct timeout_node*)t_id; + if (node->ua != ua) return ERR_FAIL; + + // Simply nullify callback - will be skipped in process_timeouts + node->callback = NULL; + + return ERR_OK; +} + + + // Instance version -err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) { - if (!ua || !t_id || !ua->timeout_heap) { - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: invalid parameters ua=%p, t_id=%p, heap=%p", - ua, t_id, ua ? ua->timeout_heap : NULL); - return ERR_FAIL; - } - - struct timeout_node* node = (struct timeout_node*)t_id; - - if (node->heap_index == SIZE_MAX || node->heap_index >= ua->timeout_heap->size || - ua->timeout_heap->heap[node->heap_index].data != node) { - DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: not found in heap: ua=%p, t_id=%p, node=%p, expires=%llu ms", - ua, t_id, node, (unsigned long long)node->expiration_ms); - return ERR_FAIL; - } - - if (timeout_heap_cancel_at(ua->timeout_heap, node->heap_index, node) == 0) { - node->heap_index = SIZE_MAX; - node->callback = NULL; - return ERR_OK; - } - - return ERR_FAIL; +err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) { + if (!ua || !t_id || !ua->timeout_heap) { + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: invalid parameters ua=%p, t_id=%p, heap=%p", + ua, t_id, ua ? ua->timeout_heap : NULL); + return ERR_FAIL; + } + + struct timeout_node* node = (struct timeout_node*)t_id; + + if (node->heap_index == SIZE_MAX || node->heap_index >= ua->timeout_heap->size || + ua->timeout_heap->heap[node->heap_index].data != node) { + DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: not found in heap: ua=%p, t_id=%p, node=%p, expires=%llu ms", + ua, t_id, node, (unsigned long long)node->expiration_ms); + return ERR_FAIL; + } + + if (timeout_heap_cancel_at(ua->timeout_heap, node->heap_index, node) == 0) { + node->heap_index = SIZE_MAX; + node->callback = NULL; + return ERR_OK; + } + + return ERR_FAIL; } @@ -657,17 +657,17 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s ua->poll_fds_dirty = 1; // Mark poll_fds as needing rebuild #if HAS_EPOLL - // Add to epoll if using epoll - if (ua->use_epoll && ua->epoll_fd >= 0) { - struct epoll_event ev; - ev.events = 0; - if (read_cbk) ev.events |= EPOLLIN; - if (write_cbk) ev.events |= EPOLLOUT; - // Embed gen in upper 32 bits of data.u64 for stale-event detection + // Add to epoll if using epoll + if (ua->use_epoll && ua->epoll_fd >= 0) { + struct epoll_event ev; + ev.events = 0; + if (read_cbk) ev.events |= EPOLLIN; + if (write_cbk) ev.events |= EPOLLOUT; + // Embed gen in upper 32 bits of data.u64 for stale-event detection ev.data.u64 = ((uint64_t)ua->sockets->sockets[index].gen << 32) | (uint32_t)fd; - - if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { - // Failed to add to epoll - remove from socket array and return error + + if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { + // Failed to add to epoll - remove from socket array and return error socket_array_remove(ua->sockets, fd); ua->socket_alloc_count--; return NULL; @@ -721,13 +721,13 @@ void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t ev.events = 0; if (read_cbk) ev.events |= EPOLLIN; if (write_cbk) ev.events |= EPOLLOUT; - // On Windows, need to cast socket_t to int for epoll_ctl -#ifdef _WIN32 - int fd = (int)(intptr_t)sock; -#else - int fd = sock; -#endif - ev.data.u64 = ((uint64_t)ua->sockets->sockets[index].gen << 32) | (uint32_t)fd; + // On Windows, need to cast socket_t to int for epoll_ctl +#ifdef _WIN32 + int fd = (int)(intptr_t)sock; +#else + int fd = sock; +#endif + ev.data.u64 = ((uint64_t)ua->sockets->sockets[index].gen << 32) | (uint32_t)fd; if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { socket_array_remove(ua->sockets, fd); ua->socket_alloc_count--; @@ -768,83 +768,83 @@ err_t uasync_remove_socket_t(struct UASYNC* ua, socket_t sock) { ua->poll_fds_dirty = 1; return ERR_OK; } - return ERR_FAIL; -} - -err_t uasync_set_socket_read(struct UASYNC* ua, void* s_id, int enable) { - if (!ua || !s_id) return ERR_FAIL; - - struct socket_node* node = (struct socket_node*)s_id; - if (!node->active || node->fd < 0) return ERR_FAIL; - - int val = enable ? 1 : 0; - if (node->enable_read == val) return ERR_OK; - node->enable_read = val; - -#if HAS_EPOLL - if (ua->use_epoll && ua->epoll_fd >= 0) { - struct epoll_event ev; - ev.events = 0; - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->read_cbk_sock && node->enable_read) ev.events |= EPOLLIN; - if (node->write_cbk_sock && node->enable_write) ev.events |= EPOLLOUT; - ev.data.fd = node->sock; - } else { - if (node->read_cbk && node->enable_read) ev.events |= EPOLLIN; - if (node->write_cbk && node->enable_write) ev.events |= EPOLLOUT; - } - if (node->except_cbk) ev.events |= EPOLLPRI; - ev.data.u64 = ((uint64_t)node->gen << 32) | (uint32_t)(node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); -#ifdef _WIN32 - int efd = (int)(intptr_t)(node->type == SOCKET_NODE_TYPE_SOCK ? node->sock : node->fd); -#else - int efd = (node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); -#endif - epoll_ctl(ua->epoll_fd, EPOLL_CTL_MOD, efd, &ev); - } -#endif - - ua->poll_fds_dirty = 1; - return ERR_OK; -} - -err_t uasync_set_socket_write(struct UASYNC* ua, void* s_id, int enable) { - if (!ua || !s_id) return ERR_FAIL; - - struct socket_node* node = (struct socket_node*)s_id; - if (!node->active || node->fd < 0) return ERR_FAIL; - - int val = enable ? 1 : 0; - if (node->enable_write == val) return ERR_OK; - node->enable_write = val; - -#if HAS_EPOLL - if (ua->use_epoll && ua->epoll_fd >= 0) { - struct epoll_event ev; - ev.events = 0; - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->read_cbk_sock && node->enable_read) ev.events |= EPOLLIN; - if (node->write_cbk_sock && node->enable_write) ev.events |= EPOLLOUT; - } else { - if (node->read_cbk && node->enable_read) ev.events |= EPOLLIN; - if (node->write_cbk && node->enable_write) ev.events |= EPOLLOUT; - } - if (node->except_cbk) ev.events |= EPOLLPRI; - ev.data.u64 = ((uint64_t)node->gen << 32) | (uint32_t)(node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); -#ifdef _WIN32 - int efd = (int)(intptr_t)(node->type == SOCKET_NODE_TYPE_SOCK ? node->sock : node->fd); -#else - int efd = (node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); -#endif - epoll_ctl(ua->epoll_fd, EPOLL_CTL_MOD, efd, &ev); - } -#endif - - ua->poll_fds_dirty = 1; - return ERR_OK; -} - -// Helper function to rebuild cached pollfd array + return ERR_FAIL; +} + +err_t uasync_set_socket_read(struct UASYNC* ua, void* s_id, int enable) { + if (!ua || !s_id) return ERR_FAIL; + + struct socket_node* node = (struct socket_node*)s_id; + if (!node->active || node->fd < 0) return ERR_FAIL; + + int val = enable ? 1 : 0; + if (node->enable_read == val) return ERR_OK; + node->enable_read = val; + +#if HAS_EPOLL + if (ua->use_epoll && ua->epoll_fd >= 0) { + struct epoll_event ev; + ev.events = 0; + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->read_cbk_sock && node->enable_read) ev.events |= EPOLLIN; + if (node->write_cbk_sock && node->enable_write) ev.events |= EPOLLOUT; + ev.data.fd = node->sock; + } else { + if (node->read_cbk && node->enable_read) ev.events |= EPOLLIN; + if (node->write_cbk && node->enable_write) ev.events |= EPOLLOUT; + } + if (node->except_cbk) ev.events |= EPOLLPRI; + ev.data.u64 = ((uint64_t)node->gen << 32) | (uint32_t)(node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); +#ifdef _WIN32 + int efd = (int)(intptr_t)(node->type == SOCKET_NODE_TYPE_SOCK ? node->sock : node->fd); +#else + int efd = (node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); +#endif + epoll_ctl(ua->epoll_fd, EPOLL_CTL_MOD, efd, &ev); + } +#endif + + ua->poll_fds_dirty = 1; + return ERR_OK; +} + +err_t uasync_set_socket_write(struct UASYNC* ua, void* s_id, int enable) { + if (!ua || !s_id) return ERR_FAIL; + + struct socket_node* node = (struct socket_node*)s_id; + if (!node->active || node->fd < 0) return ERR_FAIL; + + int val = enable ? 1 : 0; + if (node->enable_write == val) return ERR_OK; + node->enable_write = val; + +#if HAS_EPOLL + if (ua->use_epoll && ua->epoll_fd >= 0) { + struct epoll_event ev; + ev.events = 0; + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->read_cbk_sock && node->enable_read) ev.events |= EPOLLIN; + if (node->write_cbk_sock && node->enable_write) ev.events |= EPOLLOUT; + } else { + if (node->read_cbk && node->enable_read) ev.events |= EPOLLIN; + if (node->write_cbk && node->enable_write) ev.events |= EPOLLOUT; + } + if (node->except_cbk) ev.events |= EPOLLPRI; + ev.data.u64 = ((uint64_t)node->gen << 32) | (uint32_t)(node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); +#ifdef _WIN32 + int efd = (int)(intptr_t)(node->type == SOCKET_NODE_TYPE_SOCK ? node->sock : node->fd); +#else + int efd = (node->type == SOCKET_NODE_TYPE_SOCK ? (int)node->sock : node->fd); +#endif + epoll_ctl(ua->epoll_fd, EPOLL_CTL_MOD, efd, &ev); + } +#endif + + ua->poll_fds_dirty = 1; + return ERR_OK; +} + +// Helper function to rebuild cached pollfd array static void rebuild_poll_fds(struct UASYNC* ua) { if (!ua || !ua->sockets) return; @@ -894,13 +894,13 @@ static void rebuild_poll_fds(struct UASYNC* ua) { ua->poll_fds[idx].events = 0; ua->poll_fds[idx].revents = 0; - if (cur->type == SOCKET_NODE_TYPE_SOCK) { - if (cur->read_cbk_sock && cur->enable_read) ua->poll_fds[idx].events |= POLLIN; - if (cur->write_cbk_sock && cur->enable_write) ua->poll_fds[idx].events |= POLLOUT; - } else { - if (cur->read_cbk && cur->enable_read) ua->poll_fds[idx].events |= POLLIN; - if (cur->write_cbk && cur->enable_write) ua->poll_fds[idx].events |= POLLOUT; - } + if (cur->type == SOCKET_NODE_TYPE_SOCK) { + if (cur->read_cbk_sock && cur->enable_read) ua->poll_fds[idx].events |= POLLIN; + if (cur->write_cbk_sock && cur->enable_write) ua->poll_fds[idx].events |= POLLOUT; + } else { + if (cur->read_cbk && cur->enable_read) ua->poll_fds[idx].events |= POLLIN; + if (cur->write_cbk && cur->enable_write) ua->poll_fds[idx].events |= POLLOUT; + } if (cur->write_cbk && cur->enable_write) ua->poll_fds[idx].events |= POLLOUT; if (cur->except_cbk) ua->poll_fds[idx].events |= POLLPRI; @@ -913,61 +913,63 @@ static void rebuild_poll_fds(struct UASYNC* ua) { // Process events from epoll (Linux only) #if HAS_EPOLL -static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events, int n_events) { - for (int i = 0; i < n_events; i++) { - // Wakeup fd detected by data.fd < 0 (lower 32 bits of data.u64 = -1) - if (events[i].data.fd < 0) { - if (events[i].events & EPOLLIN) { drain_wakeup_pipe(ua); } - continue; - } - - int fd = (int)(events[i].data.u64 & 0xFFFFFFFF); - uint16_t ev_gen = (uint16_t)(events[i].data.u64 >> 32); - - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node || !node->active) continue; - if (node->gen != ev_gen) { DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "STALE epoll fd=%d ev_gen=%u node_gen=%u", fd, ev_gen, node->gen); continue; } - int local_fd = node->fd; - socket_t local_sock = node->sock; - int local_type = node->type; - void* local_ud = node->user_data; - socket_callback_t local_except = node->except_cbk; - socket_callback_t local_read = node->read_cbk; - socket_callback_t local_write = node->write_cbk; - socket_t_callback_t local_read_sock = node->read_cbk_sock; - socket_t_callback_t local_write_sock = node->write_cbk_sock; - - /* Check for error conditions first */ - if (events[i].events & (EPOLLERR | EPOLLHUP)) { - if (local_except) { - local_except(local_fd, local_ud); - } - } - - /* Read readiness - use appropriate callback based on socket type */ - if (events[i].events & EPOLLIN) { - if (local_type == SOCKET_NODE_TYPE_SOCK) { - if (local_read_sock) { - local_read_sock(local_sock, local_ud); - } - } else { - if (local_read) { - local_read(local_fd, local_ud); - } - } - } - - /* Write readiness - use appropriate callback based on socket type */ - if (events[i].events & EPOLLOUT) { - if (local_type == SOCKET_NODE_TYPE_SOCK) { - if (local_write_sock) { - local_write_sock(local_sock, local_ud); - } - } else { - if (local_write) { - local_write(local_fd, local_ud); - } - } +static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events, int n_events) { + for (int i = 0; i < n_events; i++) { + // Wakeup fd detected by data.fd < 0 (lower 32 bits of data.u64 = -1) + if (events[i].data.fd < 0) { + if (events[i].events & EPOLLIN) { drain_wakeup_pipe(ua); } + continue; + } + + int fd = (int)(events[i].data.u64 & 0xFFFFFFFF); + uint16_t ev_gen = (uint16_t)(events[i].data.u64 >> 32); + + if (fd == 0) fprintf(stderr, "DIAG EPOLL fd=0 ev=0x%x ev_gen=%u\n", events[i].events, ev_gen); + + struct socket_node* node = socket_array_get(ua->sockets, fd); + if (!node || !node->active) continue; + if (node->gen != ev_gen) { DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "EPOLL STALE fd=%d ev_gen=%u node_gen=%u — skipped", fd, ev_gen, node->gen); continue; } + int local_fd = node->fd; + socket_t local_sock = node->sock; + int local_type = node->type; + void* local_ud = node->user_data; + socket_callback_t local_except = node->except_cbk; + socket_callback_t local_read = node->read_cbk; + socket_callback_t local_write = node->write_cbk; + socket_t_callback_t local_read_sock = node->read_cbk_sock; + socket_t_callback_t local_write_sock = node->write_cbk_sock; + + /* Check for error conditions first */ + if (events[i].events & (EPOLLERR | EPOLLHUP)) { + if (local_except) { + local_except(local_fd, local_ud); + } + } + + /* Read readiness - use appropriate callback based on socket type */ + if (events[i].events & EPOLLIN) { + if (local_type == SOCKET_NODE_TYPE_SOCK) { + if (local_read_sock) { + local_read_sock(local_sock, local_ud); + } + } else { + if (local_read) { + local_read(local_fd, local_ud); + } + } + } + + /* Write readiness - use appropriate callback based on socket type */ + if (events[i].events & EPOLLOUT) { + if (local_type == SOCKET_NODE_TYPE_SOCK) { + if (local_write_sock) { + local_write_sock(local_sock, local_ud); + } + } else { + if (local_write) { + local_write(local_fd, local_ud); + } + } } } } @@ -978,10 +980,12 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { if (!ua) return; if (!ua->sockets || !ua->timeout_heap) return; - DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "poll(%d sockets, %zu timers, timeout=%d.%dms)", - ua->sockets->count, ua->timeout_heap->size, timeout_tb >= 0 ? timeout_tb / 10000 : -1, + DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "poll(%d sockets, %zu timers, timeout=%d.%dms)", + ua->sockets->count, ua->timeout_heap->size, timeout_tb >= 0 ? timeout_tb / 10000 : -1, timeout_tb >= 0 ? (timeout_tb % 10000) / 10 : 0); + fprintf(stderr, "DIAG uasync_poll sockets=%d timers=%zu\n", ua->sockets->count, ua->timeout_heap->size); + // Handle negative or zero timeout if (timeout_tb < 0) timeout_tb = -1; // Infinite wait else if (timeout_tb == 0) timeout_tb = 0; // No wait @@ -997,18 +1001,18 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { req_timeout.tv_usec = (timeout_tb % 10000) * 100; } - struct timeval poll_timeout; - if (timeout_tb < 0) { - poll_timeout = next_timeout; - } else { - if (next_timeout.tv_sec < req_timeout.tv_sec || - (next_timeout.tv_sec == req_timeout.tv_sec && next_timeout.tv_usec < req_timeout.tv_usec)) { - poll_timeout = next_timeout; - } else { - poll_timeout = req_timeout; - } - } - if (poll_timeout.tv_sec == 0 && poll_timeout.tv_usec == 0 && timeout_tb > 0) poll_timeout = req_timeout; + struct timeval poll_timeout; + if (timeout_tb < 0) { + poll_timeout = next_timeout; + } else { + if (next_timeout.tv_sec < req_timeout.tv_sec || + (next_timeout.tv_sec == req_timeout.tv_sec && next_timeout.tv_usec < req_timeout.tv_usec)) { + poll_timeout = next_timeout; + } else { + poll_timeout = req_timeout; + } + } + if (poll_timeout.tv_sec == 0 && poll_timeout.tv_usec == 0 && timeout_tb > 0) poll_timeout = req_timeout; int timeout_ms; @@ -1070,7 +1074,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { #ifdef _WIN32 Sleep(timeout_ms); #else - struct timespec ts = { timeout_ms / 1000, (timeout_ms % 1000) * 1000000 }; + struct timespec ts = { timeout_ms / 1000, (timeout_ms % 1000) * 1000000 }; nanosleep(&ts, NULL); #endif } @@ -1100,12 +1104,12 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { s = (SOCKET)node->fd; } - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->read_cbk_sock && node->enable_read) FD_SET(s, &read_fds); - if (node->write_cbk_sock && node->enable_write) FD_SET(s, &write_fds); - } else { - if (node->read_cbk && node->enable_read) FD_SET(s, &read_fds); - if (node->write_cbk && node->enable_write) FD_SET(s, &write_fds); + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->read_cbk_sock && node->enable_read) FD_SET(s, &read_fds); + if (node->write_cbk_sock && node->enable_write) FD_SET(s, &write_fds); + } else { + if (node->read_cbk && node->enable_read) FD_SET(s, &read_fds); + if (node->write_cbk && node->enable_write) FD_SET(s, &write_fds); } if (node->except_cbk) FD_SET(s, &except_fds); @@ -1143,7 +1147,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { int has_write = FD_ISSET(s, &write_fds); int has_except = FD_ISSET(s, &except_fds); - if (!has_read && !has_write && !has_except) continue; + if (!has_read && !has_write && !has_except) continue; DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "select→fd=%d r=%d w=%d e=%d", (int)s, has_read, has_write, has_except); if (has_except) { @@ -1202,7 +1206,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { /* Process socket events first to give sockets higher priority */ if (ret > 0) { for (int i = 0; i < ua->poll_fds_count; i++) { - if (ua->poll_fds[i].revents == 0) continue; + if (ua->poll_fds[i].revents == 0) continue; DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "poll→fd=%d rev=0x%x", ua->poll_fds[i].fd, ua->poll_fds[i].revents); /* Handle wakeup fd separately */ @@ -1309,15 +1313,15 @@ struct UASYNC* uasync_create(void) { ua->poll_fds_count = 0; ua->poll_fds_dirty = 1; - ua->wakeup_pipe[0] = -1; - ua->wakeup_pipe[1] = -1; - ua->wakeup_initialized = 0; - ua->posted_tasks_head = NULL; - ua->immediate_queue_head = NULL; - ua->immediate_queue_tail = NULL; - - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating SA..."); - ua->sockets = socket_array_create(16); + ua->wakeup_pipe[0] = -1; + ua->wakeup_pipe[1] = -1; + ua->wakeup_initialized = 0; + ua->posted_tasks_head = NULL; + ua->immediate_queue_head = NULL; + ua->immediate_queue_tail = NULL; + + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating SA..."); + ua->sockets = socket_array_create(16); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating SA1..."); if (!ua->sockets) { if (ua->wakeup_initialized) { @@ -1333,44 +1337,44 @@ struct UASYNC* uasync_create(void) { return NULL; } - ua->timeout_heap = timeout_heap_create(16); - if (!ua->timeout_heap) { - socket_array_destroy(ua->sockets); - if (ua->wakeup_initialized) { -#ifdef _WIN32 - closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[0]); - closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[1]); -#else - close(ua->wakeup_pipe[0]); - close(ua->wakeup_pipe[1]); -#endif - } - u_free(ua); - return NULL; - } - - // Initialize timeout pool - ua->timeout_pool = memory_pool_init(sizeof(struct timeout_node), "timeout_pool"); - if (!ua->timeout_pool) { - timeout_heap_destroy(ua->timeout_heap); - socket_array_destroy(ua->sockets); - if (ua->wakeup_initialized) { -#ifdef _WIN32 - closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[0]); - closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[1]); -#else - close(ua->wakeup_pipe[0]); - close(ua->wakeup_pipe[1]); -#endif - } - u_free(ua); - return NULL; - } - + ua->timeout_heap = timeout_heap_create(16); + if (!ua->timeout_heap) { + socket_array_destroy(ua->sockets); + if (ua->wakeup_initialized) { +#ifdef _WIN32 + closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[0]); + closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[1]); +#else + close(ua->wakeup_pipe[0]); + close(ua->wakeup_pipe[1]); +#endif + } + u_free(ua); + return NULL; + } + + // Initialize timeout pool + ua->timeout_pool = memory_pool_init(sizeof(struct timeout_node), "timeout_pool"); + if (!ua->timeout_pool) { + timeout_heap_destroy(ua->timeout_heap); + socket_array_destroy(ua->sockets); + if (ua->wakeup_initialized) { +#ifdef _WIN32 + closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[0]); + closesocket((SOCKET)(intptr_t)ua->wakeup_pipe[1]); +#else + close(ua->wakeup_pipe[0]); + close(ua->wakeup_pipe[1]); +#endif + } + u_free(ua); + return NULL; + } + // Set callback to u_free timeout nodes and update counters timeout_heap_set_free_callback(ua->timeout_heap, ua, timeout_node_free_callback); - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH1..."); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH1..."); // Initialize epoll on Linux ua->epoll_fd = -1; ua->use_epoll = 0; @@ -1439,19 +1443,19 @@ struct UASYNC* uasync_create(void) { u_free(ua); return NULL; } - ua->wakeup_initialized = 1; - -#if HAS_EPOLL - if (ua->use_epoll && ua->epoll_fd >= 0) { - struct epoll_event ev; - ev.events = EPOLLIN; - ev.data.fd = -1; - if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, ua->wakeup_pipe[0], &ev) < 0) { - DEBUG_WARN(DEBUG_CATEGORY_UASYNC, "Failed to add wakeup pipe to epoll: %s", strerror(errno)); - } - } -#endif - + ua->wakeup_initialized = 1; + +#if HAS_EPOLL + if (ua->use_epoll && ua->epoll_fd >= 0) { + struct epoll_event ev; + ev.events = EPOLLIN; + ev.data.fd = -1; + if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, ua->wakeup_pipe[0], &ev) < 0) { + DEBUG_WARN(DEBUG_CATEGORY_UASYNC, "Failed to add wakeup pipe to epoll: %s", strerror(errno)); + } + } +#endif + fcntl(ua->wakeup_pipe[0], F_SETFL, fcntl(ua->wakeup_pipe[0], F_GETFL, 0) | O_NONBLOCK); fcntl(ua->wakeup_pipe[1], F_SETFL, fcntl(ua->wakeup_pipe[1], F_GETFL, 0) | O_NONBLOCK); @@ -1466,124 +1470,124 @@ struct UASYNC* uasync_create(void) { return ua; } -// Print all resources for debugging -void uasync_print_resources(struct UASYNC* ua, const char* prefix) { - if (!ua) { - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: NULL uasync instance", prefix); - return; - } - - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: UASYNC Resource Report for %p", prefix, ua); - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer Statistics: allocated=%zu, u_freed=%zu, active=%zd", - ua->timer_alloc_count, ua->timer_free_count, - (ssize_t)(ua->timer_alloc_count - ua->timer_free_count)); - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket Statistics: allocated=%zu, u_freed=%zu, active=%zd", - ua->socket_alloc_count, ua->socket_free_count, - (ssize_t)(ua->socket_alloc_count - ua->socket_free_count)); - - // Показать активные таймеры - if (ua->timeout_heap) { - struct timeval now_tv; - get_current_time(&now_tv); - uint64_t now_ms = timeval_to_ms(&now_tv); - size_t active_timers = 0; - size_t deleted_timers = 0; - for (size_t i = 0; i < ua->timeout_heap->size; i++) { - uint64_t exp = ua->timeout_heap->heap[i].expiration; - int64_t remain = (int64_t)(exp - now_ms); - if (ua->timeout_heap->heap[i].deleted) { - deleted_timers++; - struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer DELETED: node=%p name='%s' exp=%llu remain=%lldms", - node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); - } else { - active_timers++; - struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer ACTIVE: node=%p name='%s' exp=%llu remain=%lldms", - node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); - } - } - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu, deleted: %zu, total: %zu", - active_timers, deleted_timers, ua->timeout_heap->size); - } - - // Показать активные сокеты - if (ua->sockets) { - int active_sockets = 0; - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket array capacity: %d, active: %d", - ua->sockets->capacity, ua->sockets->count); - for (int i = 0; i < ua->sockets->capacity; i++) { - if (ua->sockets->sockets[i].active) { - active_sockets++; - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket: fd=%d, active=%d", - ua->sockets->sockets[i].fd, - ua->sockets->sockets[i].active); - } - } - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Total active sockets: %d", active_sockets); - } - - DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: End of resource report", prefix); +// Print all resources for debugging +void uasync_print_resources(struct UASYNC* ua, const char* prefix) { + if (!ua) { + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: NULL uasync instance", prefix); + return; + } + + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: UASYNC Resource Report for %p", prefix, ua); + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer Statistics: allocated=%zu, u_freed=%zu, active=%zd", + ua->timer_alloc_count, ua->timer_free_count, + (ssize_t)(ua->timer_alloc_count - ua->timer_free_count)); + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket Statistics: allocated=%zu, u_freed=%zu, active=%zd", + ua->socket_alloc_count, ua->socket_free_count, + (ssize_t)(ua->socket_alloc_count - ua->socket_free_count)); + + // Показать активные таймеры + if (ua->timeout_heap) { + struct timeval now_tv; + get_current_time(&now_tv); + uint64_t now_ms = timeval_to_ms(&now_tv); + size_t active_timers = 0; + size_t deleted_timers = 0; + for (size_t i = 0; i < ua->timeout_heap->size; i++) { + uint64_t exp = ua->timeout_heap->heap[i].expiration; + int64_t remain = (int64_t)(exp - now_ms); + if (ua->timeout_heap->heap[i].deleted) { + deleted_timers++; + struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer DELETED: node=%p name='%s' exp=%llu remain=%lldms", + node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); + } else { + active_timers++; + struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer ACTIVE: node=%p name='%s' exp=%llu remain=%lldms", + node, node ? node->name : "null", (unsigned long long)exp, (long long)remain); + } + } + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu, deleted: %zu, total: %zu", + active_timers, deleted_timers, ua->timeout_heap->size); + } + + // Показать активные сокеты + if (ua->sockets) { + int active_sockets = 0; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket array capacity: %d, active: %d", + ua->sockets->capacity, ua->sockets->count); + for (int i = 0; i < ua->sockets->capacity; i++) { + if (ua->sockets->sockets[i].active) { + active_sockets++; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket: fd=%d, active=%d", + ua->sockets->sockets[i].fd, + ua->sockets->sockets[i].active); + } + } + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Total active sockets: %d", active_sockets); + } + + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: End of resource report", prefix); } // Modified function in u_async.c: uasync_destroy // Changes: Close wakeup sockets properly on Windows. -void uasync_destroy(struct UASYNC* ua, int close_fds) { - if (!ua) return; - - DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_destroy: starting cleanup for ua=%p", ua); - - // Диагностика ресурсов перед очисткой - uasync_print_resources(ua, "BEFORE_DESTROY"); - - // Check for potential memory leaks - if (ua->timer_alloc_count != ua->timer_free_count || ua->socket_alloc_count != ua->socket_free_count) { - DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "Memory leaks detected before cleanup: timers %zu/%zu, sockets %zu/%zu", - ua->timer_alloc_count, ua->timer_free_count, ua->socket_alloc_count, ua->socket_free_count); - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "Timer leak: allocated=%zu, u_freed=%zu, diff=%zd", - ua->timer_alloc_count, ua->timer_free_count, - (ssize_t)(ua->timer_alloc_count - ua->timer_free_count)); - DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "Socket leak: allocated=%zu, u_freed=%zu, diff=%zd", - ua->socket_alloc_count, ua->socket_free_count, - (ssize_t)(ua->socket_alloc_count - ua->socket_free_count)); - // Continue cleanup, will abort after if leaks remain - } - - // Free all remaining timeouts - - // Очистить immediate_queue - while (ua->immediate_queue_head) { - struct timeout_node* node = ua->immediate_queue_head; - ua->immediate_queue_head = node->next; - if (node) { - node->ua->timer_free_count++; - memory_pool_free(ua->timeout_pool, node); - } - } - ua->immediate_queue_tail = NULL; - - // Очистить heap - if (ua->timeout_heap) { - while (1) { - TimeoutEntry entry; - if (timeout_heap_pop(ua->timeout_heap, &entry) != 0) break; - struct timeout_node* node = (struct timeout_node*)entry.data; - - // Free all timer nodes (avoid double-u_free bug) - if (node) { - ua->timer_free_count++; - memory_pool_free(ua->timeout_pool, node); - } - } - timeout_heap_destroy(ua->timeout_heap); - ua->timeout_heap = NULL; - } - - // Destroy timeout pool - if (ua->timeout_pool) { - memory_pool_destroy(ua->timeout_pool); - ua->timeout_pool = NULL; +void uasync_destroy(struct UASYNC* ua, int close_fds) { + if (!ua) return; + + DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_destroy: starting cleanup for ua=%p", ua); + + // Диагностика ресурсов перед очисткой + uasync_print_resources(ua, "BEFORE_DESTROY"); + + // Check for potential memory leaks + if (ua->timer_alloc_count != ua->timer_free_count || ua->socket_alloc_count != ua->socket_free_count) { + DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "Memory leaks detected before cleanup: timers %zu/%zu, sockets %zu/%zu", + ua->timer_alloc_count, ua->timer_free_count, ua->socket_alloc_count, ua->socket_free_count); + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "Timer leak: allocated=%zu, u_freed=%zu, diff=%zd", + ua->timer_alloc_count, ua->timer_free_count, + (ssize_t)(ua->timer_alloc_count - ua->timer_free_count)); + DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "Socket leak: allocated=%zu, u_freed=%zu, diff=%zd", + ua->socket_alloc_count, ua->socket_free_count, + (ssize_t)(ua->socket_alloc_count - ua->socket_free_count)); + // Continue cleanup, will abort after if leaks remain + } + + // Free all remaining timeouts + + // Очистить immediate_queue + while (ua->immediate_queue_head) { + struct timeout_node* node = ua->immediate_queue_head; + ua->immediate_queue_head = node->next; + if (node) { + node->ua->timer_free_count++; + memory_pool_free(ua->timeout_pool, node); + } + } + ua->immediate_queue_tail = NULL; + + // Очистить heap + if (ua->timeout_heap) { + while (1) { + TimeoutEntry entry; + if (timeout_heap_pop(ua->timeout_heap, &entry) != 0) break; + struct timeout_node* node = (struct timeout_node*)entry.data; + + // Free all timer nodes (avoid double-u_free bug) + if (node) { + ua->timer_free_count++; + memory_pool_free(ua->timeout_pool, node); + } + } + timeout_heap_destroy(ua->timeout_heap); + ua->timeout_heap = NULL; + } + + // Destroy timeout pool + if (ua->timeout_pool) { + memory_pool_destroy(ua->timeout_pool); + ua->timeout_pool = NULL; } // Free all socket nodes using array approach @@ -1752,12 +1756,12 @@ int uasync_lookup_socket(struct UASYNC* ua, int fd, void** socket_id) { return (*socket_id != NULL) ? 0 : -1; } -void uasync_stop(struct UASYNC* ua) { - if (ua) ua->stop = 1; -} - -void uasync_mainloop(struct UASYNC* ua) { - while (!ua->stop) { - uasync_poll(ua, -1); - } +void uasync_stop(struct UASYNC* ua) { + if (ua) ua->stop = 1; +} + +void uasync_mainloop(struct UASYNC* ua) { + while (!ua->stop) { + uasync_poll(ua, -1); + } } diff --git a/tests/test_uasync_socket_race.c b/tests/test_uasync_socket_race.c index 0cceecab..9799014c 100644 --- a/tests/test_uasync_socket_race.c +++ b/tests/test_uasync_socket_race.c @@ -1,5 +1,6 @@ // test_uasync_socket_race.c — тест гонки fd-reuse в epoll при быстром accept/close/accept // 4 дочерних процесса, каждый 200 connect+send+close. Сервер на uasync+tcp_io. +// Проверяет что нет use-after-free при двойном epoll-событии (ERROR→free tc, WRITE→stale). #include #include #include @@ -104,9 +105,11 @@ static void test_timeout(void* arg) { int main(void) { printf("=== test_uasync_socket_race ===\n"); fflush(stdout); - debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + debug_config_init(); + for (int i = 1; i < DEBUG_CATEGORY_COUNT; i++) debug_set_category_level(i, DEBUG_LEVEL_ERROR); srand((unsigned)getpid()); - close(0); open("/dev/null", O_RDONLY); + + close(0); open("/dev/null", O_RDONLY); // резервируем fd 0 g_port = 25000 + (rand() % 10000); g_ua = uasync_create(); @@ -122,7 +125,6 @@ int main(void) { if (bind(lsock, (struct sockaddr*)&la, sizeof(la)) < 0) { printf("[FAIL] bind: %s\n", strerror(errno)); goto cleanup; } if (listen(lsock, 32) < 0) { printf("[FAIL] listen: %s\n", strerror(errno)); goto cleanup; } uasync_add_socket(g_ua, lsock, on_accept_cb, NULL, NULL, NULL); - printf(" listen fd=%d epoll ok\n", lsock); fflush(stdout); for (int i = 0; i < CHILDREN; i++) { pid_t pid = fork();