Browse Source

update timeout_heap and u_async: O(1) cancel, heap cancel bubble-up fix, epoll callback safety, uasync_stop, socket_array u_realloc

congestion
Evgeny 4 months ago
parent
commit
8994d7a834
  1. 46
      lib/timeout_heap.c
  2. 5
      lib/timeout_heap.h
  3. 403
      lib/u_async.c
  4. 3
      lib/u_async.h
  5. 2
      tests/bench_timeout_heap.c
  6. 2
      tests/test_u_async_comprehensive.c

46
lib/timeout_heap.c

@ -12,20 +12,20 @@
#define RIGHT_CHILD(i) (2 * (i) + 1)
TimeoutHeap *timeout_heap_create(size_t initial_capacity) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH1...");
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_INFO(DEBUG_CATEGORY_ETCP, "Creating TH2...");
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_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3...");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH3...");
h->size = 0;
h->capacity = initial_capacity;
h->freed_count = 0;
@ -50,18 +50,28 @@ void timeout_heap_destroy(TimeoutHeap *h) {
}
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 with parent
TimeoutEntry temp = h->heap[PARENT(i) - 1];
h->heap[PARENT(i) - 1] = h->heap[i - 1];
h->heap[i - 1] = temp;
swap_entries(h, PARENT(i) - 1, i - 1);
i = PARENT(i);
}
}
int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data) {
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);
@ -74,7 +84,10 @@ int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data) {
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);
@ -96,10 +109,7 @@ static void heapify_down(TimeoutHeap *h, size_t i) {
}
if (smallest == i) break;
// Swap
TimeoutEntry temp = h->heap[smallest - 1];
h->heap[smallest - 1] = h->heap[i - 1];
h->heap[i - 1] = temp;
swap_entries(h, smallest - 1, i - 1);
i = smallest;
}
}
@ -109,6 +119,7 @@ static void remove_root(TimeoutHeap *h) {
// Move last to root
h->heap[0] = h->heap[--h->size];
update_index(h, 0);
// Heapify down (1-based)
if (h->size > 0) {
@ -153,12 +164,23 @@ 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)) {
if (!h) return;
h->user_data = user_data;

5
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.

403
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
@ -130,15 +131,17 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s
if (fd >= FD_SETSIZE) return -1;
#endif
if (fd >= sa->capacity) {
// Need to resize - double the capacity
int new_capacity = sa->capacity * 2;
if (fd >= new_capacity) new_capacity = fd + 16;
if (fd >= new_capacity) new_capacity = fd + 16; // Ensure enough space
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));
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);
@ -146,11 +149,7 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s
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));
// Initialize new elements
for (int i = sa->capacity; i < new_capacity; i++) {
new_fd_to_index[i] = -1;
new_index_to_fd[i] = -1;
@ -159,11 +158,6 @@ static int socket_array_add_internal(struct socket_array* sa, int fd, socket_t s
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;
@ -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
// 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,7 +546,7 @@ 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
@ -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) {
@ -619,27 +612,17 @@ err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) {
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;
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;
}
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);
// 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;
}
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;
@ -671,6 +654,7 @@ 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) {
// Failed to add to epoll - remove from socket array and return error
socket_array_remove(ua->sockets, fd);
ua->socket_alloc_count--;
return NULL;
@ -678,15 +662,17 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s
}
#endif
return (void*)(intptr_t)fd;
// 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;
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
@ -716,18 +702,18 @@ void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t
ua->socket_alloc_count++;
ua->poll_fds_dirty = 1;
#ifdef _WIN32
int fd = (int)(intptr_t)sock;
#else
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;
// 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.fd = fd;
if (epoll_ctl(ua->epoll_fd, EPOLL_CTL_ADD, fd, &ev) < 0) {
socket_array_remove(ua->sockets, fd);
@ -737,7 +723,7 @@ void* uasync_add_socket_t(struct UASYNC* ua, socket_t sock, socket_t_callback_t
}
#endif
return (void*)(intptr_t)fd;
return &ua->sockets->sockets[index];
}
// Remove socket by socket_t
@ -839,68 +825,61 @@ 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.
// 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++) {
if (events[i].data.fd == -1) {
// 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);
}
}
}
@ -1066,14 +1045,10 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
if (!node->active) continue;
SOCKET s;
int fd;
int type = node->type;
if (type == SOCKET_NODE_TYPE_SOCK) {
if (node->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);
@ -1083,59 +1058,31 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
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 (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
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 (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) {
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 (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);
}
}
}
@ -1171,79 +1118,56 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
/* Handle wakeup fd separately */
if (wakeup_fd_present && i == 0) {
if (ua->poll_fds[i].revents & POLLIN) {
handle_wakeup(ua);
drain_wakeup_pipe(ua);
}
continue;
}
int fd = ua->poll_fds[i].fd;
/* 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 — fresh lookup per group */
/* Check for error conditions first */
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);
/* 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) {
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);
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
/* Read readiness — fresh lookup, copy cb+args before call */
/* Read readiness - use appropriate callback based on socket type */
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);
}
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 — fresh lookup, copy cb+args before call */
/* Write readiness - use appropriate callback based on socket type */
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);
}
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);
}
}
}
@ -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...");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating SA...");
ua->sockets = socket_array_create(16);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating SA1...");
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,7 +1303,7 @@ 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
@ -1387,9 +1311,8 @@ struct UASYNC* uasync_create(void) {
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;
u_free(ua);
return NULL;
}
struct sockaddr_in addr = {0};
@ -1403,7 +1326,8 @@ struct UASYNC* uasync_create(void) {
closesocket(r);
closesocket(w);
DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "Wakeup socket setup failed: %d", WSAGetLastError());
goto create_cleanup;
u_free(ua);
return NULL;
}
ua->wakeup_pipe[0] = (int)(intptr_t)r;
@ -1422,10 +1346,22 @@ struct UASYNC* uasync_create(void) {
// POSIX pipe
if (pipe(ua->wakeup_pipe) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno));
goto create_cleanup;
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
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);
@ -1435,33 +1371,23 @@ struct UASYNC* uasync_create(void) {
pthread_mutex_init(&ua->posted_lock, NULL);
#endif
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3...");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH3...");
return ua;
create_cleanup:
memory_pool_destroy(ua->timeout_pool);
timeout_heap_destroy(ua->timeout_heap);
socket_array_destroy(ua->sockets);
#if HAS_EPOLL
if (ua->epoll_fd >= 0) close(ua->epoll_fd);
#endif
u_free(ua);
return NULL;
}
// Print all resources for debugging
void uasync_print_resources(struct UASYNC* ua, const char* prefix) {
if (!ua) {
printf("%s: NULL uasync instance\n", prefix);
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: NULL uasync instance", prefix);
return;
}
printf("\n🔍 %s: UASYNC Resource Report for %p\n", prefix, ua);
printf(" Timer Statistics: allocated=%zu, u_freed=%zu, active=%zd\n",
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));
printf(" Socket Statistics: allocated=%zu, u_freed=%zu, active=%zd\n",
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));
@ -1473,30 +1399,30 @@ void uasync_print_resources(struct UASYNC* ua, const char* prefix) {
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",
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Timer: node=%p, expires=%llu ms",
node, (unsigned long long)ua->timeout_heap->heap[i].expiration);
}
}
printf(" Active timers in heap: %zu\n", active_timers);
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Active timers in heap: %zu", active_timers);
}
// Показать активные сокеты
if (ua->sockets) {
int active_sockets = 0;
printf(" Socket array capacity: %d, active: %d\n",
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++;
printf(" Socket: fd=%d, active=%d\n",
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Socket: fd=%d, active=%d",
ua->sockets->sockets[i].fd,
ua->sockets->sockets[i].active);
}
}
printf(" Total active sockets: %d\n", active_sockets);
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, " Total active sockets: %d", active_sockets);
}
printf("🔚 %s: End of resource report\n\n", prefix);
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "%s: End of resource report", prefix);
}
// Modified function in u_async.c: uasync_destroy
@ -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_stop(struct UASYNC* ua) {
if (ua) ua->stop = 1;
}
void uasync_mainloop(struct UASYNC* ua) {
while (1) {
uasync_poll(ua, -1); // Infinite wait
}
while (!ua->stop) {
uasync_poll(ua, -1);
}
}

3
lib/u_async.h

@ -7,6 +7,7 @@
#include "platform_compat.h"
#include <stddef.h>
#include <signal.h>
#include "timeout_heap.h"
#include "socket_compat.h"
@ -68,6 +69,7 @@ struct UASYNC {
#else
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);

2
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;

2
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) {

Loading…
Cancel
Save