From 8994d7a8348409b103450a193cbf7fdf4847780a Mon Sep 17 00:00:00 2001 From: Evgeny Date: Fri, 29 May 2026 01:11:37 +0300 Subject: [PATCH] update timeout_heap and u_async: O(1) cancel, heap cancel bubble-up fix, epoll callback safety, uasync_stop, socket_array u_realloc --- lib/timeout_heap.c | 184 ++++--- lib/timeout_heap.h | 5 +- lib/u_async.c | 843 +++++++++++++---------------- lib/u_async.h | 13 +- tests/bench_timeout_heap.c | 2 +- tests/test_u_async_comprehensive.c | 2 +- 6 files changed, 503 insertions(+), 546 deletions(-) diff --git a/lib/timeout_heap.c b/lib/timeout_heap.c index cc82ee87..e91cf50b 100644 --- a/lib/timeout_heap.c +++ b/lib/timeout_heap.c @@ -11,21 +11,21 @@ #define LEFT_CHILD(i) (2 * (i)) #define RIGHT_CHILD(i) (2 * (i) + 1) -TimeoutHeap *timeout_heap_create(size_t initial_capacity) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH1..."); - TimeoutHeap *h = u_malloc(sizeof(TimeoutHeap)); - if (!h) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH0 error..."); - return NULL; - } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH2..."); - h->heap = u_malloc(sizeof(TimeoutEntry) * initial_capacity); - if (!h->heap) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH1 error..."); - u_free(h); - return NULL; - } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3..."); +TimeoutHeap *timeout_heap_create(size_t initial_capacity) { + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH1..."); + TimeoutHeap *h = u_malloc(sizeof(TimeoutHeap)); + if (!h) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH0 error..."); + return NULL; + } + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH2..."); + h->heap = u_malloc(sizeof(TimeoutEntry) * initial_capacity); + if (!h->heap) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH1 error..."); + u_free(h); + return NULL; + } + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH3..."); h->size = 0; h->capacity = initial_capacity; h->freed_count = 0; @@ -50,70 +50,81 @@ void timeout_heap_destroy(TimeoutHeap *h) { } -static void bubble_up(TimeoutHeap *h, size_t i) { - // i is 1-based - while (i > 1 && h->heap[PARENT(i) - 1].expiration > h->heap[i - 1].expiration) { - // Swap with parent - TimeoutEntry temp = h->heap[PARENT(i) - 1]; - h->heap[PARENT(i) - 1] = h->heap[i - 1]; - h->heap[i - 1] = temp; - i = PARENT(i); - } +static void update_index(TimeoutHeap *h, size_t idx) { + if (h->heap[idx].index_ptr) + *h->heap[idx].index_ptr = idx; +} + +static void swap_entries(TimeoutHeap *h, size_t a, size_t b) { + TimeoutEntry temp = h->heap[a]; + h->heap[a] = h->heap[b]; + h->heap[b] = temp; + update_index(h, a); + update_index(h, b); +} + +static void bubble_up(TimeoutHeap *h, size_t i) { + // i is 1-based + while (i > 1 && h->heap[PARENT(i) - 1].expiration > h->heap[i - 1].expiration) { + swap_entries(h, PARENT(i) - 1, i - 1); + i = PARENT(i); + } } -int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data) { - if (h->size == h->capacity) { - size_t new_cap = h->capacity ? h->capacity * 2 : 1; - TimeoutEntry *new_heap = u_realloc(h->heap, sizeof(TimeoutEntry) * new_cap); - if (!new_heap) return -1; // Allocation failed - h->heap = new_heap; - h->capacity = new_cap; - } - - // Insert at end (0-based) - size_t idx = h->size++; - h->heap[idx].expiration = expiration; - h->heap[idx].data = data; - h->heap[idx].deleted = 0; - - // Bubble up (1-based) - bubble_up(h, idx + 1); - return 0; +int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data, size_t *index_ptr) { + if (h->size == h->capacity) { + size_t new_cap = h->capacity ? h->capacity * 2 : 1; + TimeoutEntry *new_heap = u_realloc(h->heap, sizeof(TimeoutEntry) * new_cap); + if (!new_heap) return -1; // Allocation failed + h->heap = new_heap; + h->capacity = new_cap; + } + + // Insert at end (0-based) + size_t idx = h->size++; + h->heap[idx].expiration = expiration; + h->heap[idx].data = data; + h->heap[idx].index_ptr = index_ptr; + h->heap[idx].deleted = 0; + if (index_ptr) + *index_ptr = idx; + + // Bubble up (1-based) + bubble_up(h, idx + 1); + return 0; } -static void heapify_down(TimeoutHeap *h, size_t i) { - // i is 1-based - while (1) { - size_t smallest = i; - size_t left = LEFT_CHILD(i); - size_t right = RIGHT_CHILD(i); - - if (left <= h->size && h->heap[left - 1].expiration < h->heap[smallest - 1].expiration) { - smallest = left; - } - if (right <= h->size && h->heap[right - 1].expiration < h->heap[smallest - 1].expiration) { - smallest = right; - } - if (smallest == i) break; - - // Swap - TimeoutEntry temp = h->heap[smallest - 1]; - h->heap[smallest - 1] = h->heap[i - 1]; - h->heap[i - 1] = temp; - i = smallest; - } +static void heapify_down(TimeoutHeap *h, size_t i) { + // i is 1-based + while (1) { + size_t smallest = i; + size_t left = LEFT_CHILD(i); + size_t right = RIGHT_CHILD(i); + + if (left <= h->size && h->heap[left - 1].expiration < h->heap[smallest - 1].expiration) { + smallest = left; + } + if (right <= h->size && h->heap[right - 1].expiration < h->heap[smallest - 1].expiration) { + smallest = right; + } + if (smallest == i) break; + + swap_entries(h, smallest - 1, i - 1); + i = smallest; + } } -static void remove_root(TimeoutHeap *h) { - if (h->size == 0) return; - - // Move last to root - h->heap[0] = h->heap[--h->size]; - - // Heapify down (1-based) - if (h->size > 0) { - heapify_down(h, 1); - } +static void remove_root(TimeoutHeap *h) { + if (h->size == 0) return; + + // Move last to root + h->heap[0] = h->heap[--h->size]; + update_index(h, 0); + + // Heapify down (1-based) + if (h->size > 0) { + heapify_down(h, 1); + } } int timeout_heap_peek(TimeoutHeap *h, TimeoutEntry *out) { @@ -149,14 +160,25 @@ int timeout_heap_pop(TimeoutHeap *h, TimeoutEntry *out) { return 0; } -int timeout_heap_cancel(TimeoutHeap *h, TimeoutTime expiration, void *data) { - for (size_t i = 0; i < h->size; ++i) { - if (h->heap[i].expiration == expiration && h->heap[i].data == data) { - h->heap[i].deleted = 1; - return 0; - } - } - return -1; // Not found +int timeout_heap_cancel(TimeoutHeap *h, TimeoutTime expiration, void *data) { + for (size_t i = 0; i < h->size; ++i) { + if (h->heap[i].expiration == expiration && h->heap[i].data == data) { + h->heap[i].deleted = 1; + h->heap[i].expiration = 0; + bubble_up(h, i + 1); + return 0; + } + } + return -1; // Not found +} + +int timeout_heap_cancel_at(TimeoutHeap *h, size_t index, void *data) { + if (index >= h->size || h->heap[index].data != data) return -1; + + h->heap[index].deleted = 1; + h->heap[index].expiration = 0; + bubble_up(h, index + 1); + return 0; } void timeout_heap_set_free_callback(TimeoutHeap *h, void* user_data, void (*callback)(void* user_data, void* data)) { diff --git a/lib/timeout_heap.h b/lib/timeout_heap.h index 8cf979f6..8e093f2f 100644 --- a/lib/timeout_heap.h +++ b/lib/timeout_heap.h @@ -11,6 +11,7 @@ typedef uint64_t TimeoutTime; // e.g., milliseconds since epoch or from now typedef struct { TimeoutTime expiration; // Sort key (smaller = earlier) void *data; // User data (e.g., callback or ID) + size_t *index_ptr; // Pointer to node's heap_index (NULL = no tracking) int deleted; // 0 = active, 1 = deleted } TimeoutEntry; @@ -53,7 +54,7 @@ void timeout_heap_set_free_callback(TimeoutHeap *h, void* user_data, void (*call * @param data User data associated with the timeout. * @return 0 on success, -1 on allocation failure. */ -int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data); +int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data, size_t *index_ptr); /** * Peek at the earliest non-deleted timeout without removing it. @@ -82,6 +83,8 @@ int timeout_heap_pop(TimeoutHeap *h, TimeoutEntry *out); */ int timeout_heap_cancel(TimeoutHeap *h, TimeoutTime expiration, void *data); +int timeout_heap_cancel_at(TimeoutHeap *h, size_t index, void *data); + /** * Get the number of freed timer nodes. * @param h The heap. diff --git a/lib/u_async.c b/lib/u_async.c index d1645eef..9dbf47dd 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -38,6 +38,7 @@ struct timeout_node { 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 @@ -129,46 +130,39 @@ 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) { - 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; + 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; } // Check if FD already exists @@ -235,11 +229,17 @@ 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 - sa->sockets[index].active = 0; - sa->sockets[index].fd = -1; - sa->sockets[index].sock = SOCKET_INVALID; - sa->sockets[index].type = SOCKET_NODE_TYPE_FD; + // 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->fd_to_index[fd] = -1; sa->index_to_fd[index] = -1; @@ -537,15 +537,7 @@ void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_c node->arg = arg; node->callback = callback; node->ua = ua; - - if (timeout_tb == 0) { - node->expiration_ms = 0; - node->next = NULL; - if (ua->immediate_queue_tail) ua->immediate_queue_tail->next = node; - else ua->immediate_queue_head = node; - ua->immediate_queue_tail = node; - return node; - } + node->heap_index = SIZE_MAX; // Calculate expiration time in milliseconds struct timeval now; @@ -554,14 +546,14 @@ void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_c node->expiration_ms = timeval_to_ms(&now); // Add to heap - if (timeout_heap_push(ua->timeout_heap, node->expiration_ms, node) != 0) { + 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; + } + + return node; } // Immediate execution in next mainloop (FIFO order) @@ -582,6 +574,7 @@ void* uasync_call_soon(struct UASYNC* ua, void* user_arg, timeout_callback_t cal node->ua = ua; node->expiration_ms = 0; node->next = NULL; + node->heap_index = SIZE_MAX; // FIFO: добавляем в конец очереди if (ua->immediate_queue_tail) { @@ -610,39 +603,29 @@ err_t uasync_call_soon_cancel(struct UASYNC* ua, void* t_id) { // 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; - - // Try to cancel from heap first - if (timeout_heap_cancel(ua->timeout_heap, node->expiration_ms, node) == 0) { - // Successfully marked as deleted - u_free will happen lazily in heap - node->callback = NULL; -// DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: successfully cancelled timer %p from heap", node); - return ERR_OK; - } - - 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); +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; + } - // Not in heap — try immediate_queue (for timeout=0 timers placed via FIFO) - { - struct timeout_node* cur = ua->immediate_queue_head; - while (cur) { - if (cur == node) { - cur->callback = NULL; - return ERR_OK; - } - cur = cur->next; - } + 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; + return ERR_FAIL; } @@ -671,22 +654,25 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s ev.data.fd = fd; if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { - socket_array_remove(ua->sockets, fd); - ua->socket_alloc_count--; - return NULL; - } - } -#endif - - return (void*)(intptr_t)fd; + // Failed to add to epoll - remove from socket array and return error + socket_array_remove(ua->sockets, fd); + ua->socket_alloc_count--; + return NULL; + } + } +#endif + + // Return pointer to the socket node as ID + return &ua->sockets->sockets[index]; } -err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) { - if (!ua || !s_id) return ERR_FAIL; - - int fd = (int)(intptr_t)s_id; - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node || !node->active) return ERR_FAIL; +err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) { + 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 fd = node->fd; #if HAS_EPOLL // Remove from epoll if using epoll @@ -713,31 +699,31 @@ void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t (socket_callback_t)except_cbk, user_data); if (index < 0) return NULL; - ua->socket_alloc_count++; - ua->poll_fds_dirty = 1; - + ua->socket_alloc_count++; + ua->poll_fds_dirty = 1; + +#if HAS_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; + // On Windows, need to cast socket_t to int for epoll_ctl #ifdef _WIN32 - int fd = (int)(intptr_t)sock; + int fd = (int)(intptr_t)sock; #else - int fd = sock; + int fd = sock; #endif - -#if HAS_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; ev.data.fd = fd; - if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { - socket_array_remove(ua->sockets, fd); - ua->socket_alloc_count--; - return NULL; - } - } -#endif - - return (void*)(intptr_t)fd; + if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) { + socket_array_remove(ua->sockets, fd); + ua->socket_alloc_count--; + return NULL; + } + } +#endif + + return &ua->sockets->sockets[index]; } // Remove socket by socket_t @@ -839,73 +825,66 @@ static void rebuild_poll_fds(struct UASYNC* ua) { ua->poll_fds_dirty = 0; } -// Process events from epoll (Linux only). Each callback group re-looks up node -// via socket_array_get to survive realloc triggered by other callbacks. -#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++) { - if (events[i].data.fd == -1) { +// 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++) { + // Check if this is the wakeup fd (data.fd is -1) + if (events[i].data.fd < 0) { if (events[i].events & EPOLLIN) { - handle_wakeup(ua); + drain_wakeup_pipe(ua); } continue; } - int fd = events[i].data.fd; + // Socket event — save node fields locally: callbacks may realloc + // the socket array, invalidating `node`. + struct socket_node* node = socket_array_get(ua->sockets, events[i].data.fd); + if (!node || !node->active) 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)) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (node && node->active && node->except_cbk) { - socket_callback_t cb = node->except_cbk; - int node_fd = node->fd; - void* ud = node->user_data; - cb(node_fd, ud); + if (local_except) { + local_except(local_fd, local_ud); } } + /* Read readiness - use appropriate callback based on socket type */ if (events[i].events & EPOLLIN) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (node && node->active) { - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->read_cbk_sock) { - socket_t_callback_t cb = node->read_cbk_sock; - socket_t sock = node->sock; - void* ud = node->user_data; - cb(sock, ud); - } - } else { - if (node->read_cbk) { - socket_callback_t cb = node->read_cbk; - int node_fd = node->fd; - void* ud = node->user_data; - cb(node_fd, ud); - } + 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) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (node && node->active) { - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->write_cbk_sock) { - socket_t_callback_t cb = node->write_cbk_sock; - socket_t sock = node->sock; - void* ud = node->user_data; - cb(sock, ud); - } - } else { - if (node->write_cbk) { - socket_callback_t cb = node->write_cbk; - int node_fd = node->fd; - void* ud = node->user_data; - cb(node_fd, ud); - } + 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); } } - } - } -} + } + } +} #endif // Instance version @@ -1059,87 +1038,55 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { return; } - if (ret > 0) { - for (int i = 0; i < ua->sockets->count; i++) { - int idx = ua->sockets->active_indices[i]; - struct socket_node* node = &ua->sockets->sockets[idx]; - if (!node->active) continue; - - SOCKET s; - int fd; - int type = node->type; - if (type == SOCKET_NODE_TYPE_SOCK) { - s = node->sock; - fd = (int)(intptr_t)s; - } else { - s = (SOCKET)node->fd; - fd = node->fd; - } - - int has_read = FD_ISSET(s, &read_fds); - 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_except) { - struct socket_node* n = (type == SOCKET_NODE_TYPE_SOCK) ? - socket_array_get_by_sock(ua->sockets, s) : - socket_array_get(ua->sockets, fd); - if (n && n->active && n->except_cbk) { - socket_callback_t cb = n->except_cbk; - int n_fd = n->fd; - void* ud = n->user_data; - cb(n_fd, ud); - } - } - - if (has_read) { - struct socket_node* n = (type == SOCKET_NODE_TYPE_SOCK) ? - socket_array_get_by_sock(ua->sockets, s) : - socket_array_get(ua->sockets, fd); - if (n && n->active) { - if (type == SOCKET_NODE_TYPE_SOCK) { - if (n->read_cbk_sock) { - socket_t_callback_t cb = n->read_cbk_sock; - socket_t sock = n->sock; - void* ud = n->user_data; - cb(sock, ud); - } - } else { - if (n->read_cbk) { - socket_callback_t cb = n->read_cbk; - int n_fd = n->fd; - void* ud = n->user_data; - cb(n_fd, ud); - } - } - } - } - - if (has_write) { - struct socket_node* n = (type == SOCKET_NODE_TYPE_SOCK) ? - socket_array_get_by_sock(ua->sockets, s) : - socket_array_get(ua->sockets, fd); - if (n && n->active) { - if (type == SOCKET_NODE_TYPE_SOCK) { - if (n->write_cbk_sock) { - socket_t_callback_t cb = n->write_cbk_sock; - socket_t sock = n->sock; - void* ud = n->user_data; - cb(sock, ud); - } - } else { - if (n->write_cbk) { - socket_callback_t cb = n->write_cbk; - int n_fd = n->fd; - void* ud = n->user_data; - cb(n_fd, ud); - } - } - } - } - } + if (ret > 0) { + for (int i = 0; i < ua->sockets->count; i++) { + int idx = ua->sockets->active_indices[i]; + struct socket_node* node = &ua->sockets->sockets[idx]; + if (!node->active) continue; + + SOCKET s; + if (node->type == SOCKET_NODE_TYPE_SOCK) { + s = node->sock; + } else { + s = (SOCKET)node->fd; + } + + int has_read = FD_ISSET(s, &read_fds); + 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_except) { + if (node->except_cbk) { + node->except_cbk(node->fd, node->user_data); + } + } + + if (has_read) { + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->read_cbk_sock) { + node->read_cbk_sock(node->sock, node->user_data); + } + } else { + if (node->read_cbk) { + node->read_cbk(node->fd, node->user_data); + } + } + } + + if (has_write) { + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->write_cbk_sock) { + node->write_cbk_sock(node->sock, node->user_data); + } + } else { + if (node->write_cbk) { + node->write_cbk(node->fd, node->user_data); + } + } + } + } } #else // On non-Windows, use poll() @@ -1163,91 +1110,68 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) { return; } - /* 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; - - /* Handle wakeup fd separately */ - if (wakeup_fd_present && i == 0) { - if (ua->poll_fds[i].revents & POLLIN) { - handle_wakeup(ua); - } - continue; - } - - int fd = ua->poll_fds[i].fd; - - /* Check for error conditions first — fresh lookup per group */ - if (ua->poll_fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node) node = socket_array_get_by_sock(ua->sockets, (socket_t)fd); - if (node && node->active && node->except_cbk) { - socket_callback_t cb = node->except_cbk; - int n_fd = node->fd; - void* ud = node->user_data; - cb(n_fd, ud); - } - } - - /* Exceptional data (out-of-band) */ - if (ua->poll_fds[i].revents & POLLPRI) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node) node = socket_array_get_by_sock(ua->sockets, (socket_t)fd); - if (node && node->active && node->except_cbk) { - socket_callback_t cb = node->except_cbk; - int n_fd = node->fd; - void* ud = node->user_data; - cb(n_fd, ud); - } - } - - /* Read readiness — fresh lookup, copy cb+args before call */ - if (ua->poll_fds[i].revents & POLLIN) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node) node = socket_array_get_by_sock(ua->sockets, (socket_t)fd); - if (node && node->active) { - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->read_cbk_sock) { - socket_t_callback_t cb = node->read_cbk_sock; - socket_t sock = node->sock; - void* ud = node->user_data; - cb(sock, ud); - } - } else { - if (node->read_cbk) { - socket_callback_t cb = node->read_cbk; - int n_fd = node->fd; - void* ud = node->user_data; - cb(n_fd, ud); - } - } - } - } - - /* Write readiness — fresh lookup, copy cb+args before call */ - if (ua->poll_fds[i].revents & POLLOUT) { - struct socket_node* node = socket_array_get(ua->sockets, fd); - if (!node) node = socket_array_get_by_sock(ua->sockets, (socket_t)fd); - if (node && node->active) { - if (node->type == SOCKET_NODE_TYPE_SOCK) { - if (node->write_cbk_sock) { - socket_t_callback_t cb = node->write_cbk_sock; - socket_t sock = node->sock; - void* ud = node->user_data; - cb(sock, ud); - } - } else { - if (node->write_cbk) { - socket_callback_t cb = node->write_cbk; - int n_fd = node->fd; - void* ud = node->user_data; - cb(n_fd, ud); - } - } - } - } - } + /* 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; + + /* Handle wakeup fd separately */ + if (wakeup_fd_present && i == 0) { + if (ua->poll_fds[i].revents & POLLIN) { + drain_wakeup_pipe(ua); + } + continue; + } + + /* Socket event - lookup by fd */ + struct socket_node* node = socket_array_get(ua->sockets, ua->poll_fds[i].fd); + if (!node) { // Try by socket_t (in case this is a socket) + socket_t lookup_sock = ua->poll_fds[i].fd; + node = socket_array_get_by_sock(ua->sockets, lookup_sock); + } + if (!node) continue; // Socket may have been removed + + /* Check for error conditions first */ + if (ua->poll_fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) { + /* Treat as exceptional condition */ + if (node->except_cbk) { + node->except_cbk(node->fd, node->user_data); + } + } + + /* Exceptional data (out-of-band) */ + if (ua->poll_fds[i].revents & POLLPRI) { + if (node->except_cbk) { + node->except_cbk(node->fd, node->user_data); + } + } + + /* Read readiness - use appropriate callback based on socket type */ + if (ua->poll_fds[i].revents & POLLIN) { + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->read_cbk_sock) { + node->read_cbk_sock(node->sock, node->user_data); + } + } else { + if (node->read_cbk) { + node->read_cbk(node->fd, node->user_data); + } + } + } + + /* Write readiness - use appropriate callback based on socket type */ + if (ua->poll_fds[i].revents & POLLOUT) { + if (node->type == SOCKET_NODE_TYPE_SOCK) { + if (node->write_cbk_sock) { + node->write_cbk_sock(node->sock, node->user_data); + } + } else { + if (node->write_cbk) { + node->write_cbk(node->fd, node->user_data); + } + } + } + } } #endif @@ -1302,9 +1226,9 @@ struct UASYNC* uasync_create(void) { ua->immediate_queue_head = NULL; ua->immediate_queue_tail = NULL; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating SA..."); - ua->sockets = socket_array_create(16); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating SA1..."); + 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) { #ifdef _WIN32 @@ -1356,7 +1280,7 @@ struct UASYNC* uasync_create(void) { // Set callback to u_free timeout nodes and update counters timeout_heap_set_free_callback(ua->timeout_heap, ua, timeout_node_free_callback); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH1..."); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH1..."); // Initialize epoll on Linux ua->epoll_fd = -1; ua->use_epoll = 0; @@ -1369,7 +1293,7 @@ struct UASYNC* uasync_create(void) { if (ua->wakeup_initialized) { struct epoll_event ev; ev.events = EPOLLIN; - ev.data.fd = -1; + ev.data.fd = -1; // -1 fd indicates wakeup 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)); } @@ -1379,31 +1303,31 @@ struct UASYNC* uasync_create(void) { } #endif - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH2..."); + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH2..."); #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"); - 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; + 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; } ua->wakeup_pipe[0] = (int)(intptr_t)r; @@ -1420,12 +1344,24 @@ struct UASYNC* uasync_create(void) { #else // POSIX pipe - if (pipe(ua->wakeup_pipe) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno)); - goto create_cleanup; + if (pipe(ua->wakeup_pipe) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno)); + u_free(ua); + return NULL; } - ua->wakeup_initialized = 1; - + 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); @@ -1434,95 +1370,85 @@ struct UASYNC* uasync_create(void) { } pthread_mutex_init(&ua->posted_lock, NULL); -#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; +#endif + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH3..."); + + return ua; } -// Print all resources for debugging -void uasync_print_resources(struct UASYNC* ua, const char* prefix) { - if (!ua) { - printf("%s: NULL uasync instance\n", prefix); - return; - } - - printf("\n🔍 %s: UASYNC Resource Report for %p\n", prefix, ua); - printf(" Timer Statistics: allocated=%zu, u_freed=%zu, active=%zd\n", - ua->timer_alloc_count, ua->timer_free_count, - (ssize_t)(ua->timer_alloc_count - ua->timer_free_count)); - printf(" Socket Statistics: allocated=%zu, u_freed=%zu, active=%zd\n", - ua->socket_alloc_count, ua->socket_free_count, - (ssize_t)(ua->socket_alloc_count - ua->socket_free_count)); - - // Показать активные таймеры - if (ua->timeout_heap) { - size_t active_timers = 0; - // Безопасное чтение без извлечения - просто итерируем по массиву - for (size_t i = 0; i < ua->timeout_heap->size; i++) { - if (!ua->timeout_heap->heap[i].deleted) { - active_timers++; - struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; - printf(" Timer: node=%p, expires=%llu ms\n", - node, (unsigned long long)ua->timeout_heap->heap[i].expiration); - } - } - printf(" Active timers in heap: %zu\n", active_timers); - } - - // Показать активные сокеты - if (ua->sockets) { - int active_sockets = 0; - printf(" Socket array capacity: %d, active: %d\n", - ua->sockets->capacity, ua->sockets->count); - for (int i = 0; i < ua->sockets->capacity; i++) { - if (ua->sockets->sockets[i].active) { - active_sockets++; - printf(" Socket: fd=%d, active=%d\n", - ua->sockets->sockets[i].fd, - ua->sockets->sockets[i].active); - } - } - printf(" Total active sockets: %d\n", active_sockets); - } - - printf("🔚 %s: End of resource report\n\n", 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) { + size_t active_timers = 0; + // Безопасное чтение без извлечения - просто итерируем по массиву + for (size_t i = 0; i < ua->timeout_heap->size; i++) { + if (!ua->timeout_heap->heap[i].deleted) { + active_timers++; + struct timeout_node* node = (struct timeout_node*)ua->timeout_heap->heap[i].data; + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer: node=%p, expires=%llu ms", + node, (unsigned long long)ua->timeout_heap->heap[i].expiration); + } + } + DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu", active_timers); + } + + // Показать активные сокеты + 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 - } - +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 @@ -1538,7 +1464,6 @@ void uasync_destroy(struct UASYNC* ua, int close_fds) { // Очистить heap if (ua->timeout_heap) { - size_t u_freed_count = 0; while (1) { TimeoutEntry entry; if (timeout_heap_pop(ua->timeout_heap, &entry) != 0) break; @@ -1726,8 +1651,12 @@ int uasync_lookup_socket(struct UASYNC* ua, int fd, void** socket_id) { return (*socket_id != NULL) ? 0 : -1; } -void uasync_mainloop(struct UASYNC* ua) { - while (1) { - uasync_poll(ua, -1); // Infinite wait - } +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/lib/u_async.h b/lib/u_async.h index a78b8f8d..fbfc30fb 100644 --- a/lib/u_async.h +++ b/lib/u_async.h @@ -5,9 +5,10 @@ #ifndef UASYNC_H #define UASYNC_H -#include "platform_compat.h" -#include -#include "timeout_heap.h" +#include "platform_compat.h" +#include +#include +#include "timeout_heap.h" #include "socket_compat.h" typedef void (*timeout_callback_t)(void* user_arg);// передаёт user_arg из uasync_set_timeout @@ -66,8 +67,9 @@ struct UASYNC { #ifdef _WIN32 CRITICAL_SECTION posted_lock; #else - pthread_mutex_t posted_lock; -#endif + pthread_mutex_t posted_lock; +#endif + volatile sig_atomic_t stop; }; // Type definitions @@ -103,6 +105,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb); // Mainloop (бесконечный цикл, __noreturn) void uasync_mainloop(struct UASYNC* ua); +void uasync_stop(struct UASYNC* ua); // Debug statistics void uasync_get_stats(struct UASYNC* ua, size_t* timer_alloc, size_t* timer_free, size_t* socket_alloc, size_t* socket_free); diff --git a/tests/bench_timeout_heap.c b/tests/bench_timeout_heap.c index 1ed76481..72e9f756 100644 --- a/tests/bench_timeout_heap.c +++ b/tests/bench_timeout_heap.c @@ -54,7 +54,7 @@ int main(void) { // Benchmark 1: 100 random pushes clock_gettime(CLOCK_MONOTONIC, &start); for (i = 0; i < NUM_OPERATIONS; i++) { - if (timeout_heap_push(heap, timeouts[i], data[i]) != 0) { + if (timeout_heap_push(heap, timeouts[i], data[i], NULL) != 0) { fprintf(stderr, "Push failed at iteration %d\n", i); timeout_heap_destroy(heap); return 1; diff --git a/tests/test_u_async_comprehensive.c b/tests/test_u_async_comprehensive.c index 292a89fb..41e587d9 100644 --- a/tests/test_u_async_comprehensive.c +++ b/tests/test_u_async_comprehensive.c @@ -267,7 +267,7 @@ static void order_immediate_cb(void* arg); static void order_delayed_cb(void* arg) { struct ordering_ctx* ctx = (struct ordering_ctx*)arg; ctx->seq[ctx->idx++] = 1; - if (ctx->idx == 1) uasync_set_timeout(ctx->ua, 0, ctx, order_immediate_cb, "test_order_imm"); + if (ctx->idx == 1) uasync_call_soon(ctx->ua, ctx, order_immediate_cb); } static void order_immediate_cb(void* arg) {