Browse Source

Optimize uasync_poll: cache pollfd array, add active_indices for O(1) socket traversal

nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
4f287c711f
  1. 160
      lib/u_async.c
  2. 31
      lib/u_async.h

160
lib/u_async.c

@ -37,6 +37,7 @@ 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
@ -60,11 +61,13 @@ static struct socket_array* socket_array_create(int initial_capacity) {
sa->sockets = calloc(initial_capacity, sizeof(struct socket_node));
sa->fd_to_index = calloc(initial_capacity, sizeof(int));
sa->index_to_fd = calloc(initial_capacity, sizeof(int));
sa->active_indices = calloc(initial_capacity, sizeof(int));
if (!sa->sockets || !sa->fd_to_index || !sa->index_to_fd) {
if (!sa->sockets || !sa->fd_to_index || !sa->index_to_fd || !sa->active_indices) {
free(sa->sockets);
free(sa->fd_to_index);
free(sa->index_to_fd);
free(sa->active_indices);
free(sa);
return NULL;
}
@ -73,6 +76,7 @@ static struct socket_array* socket_array_create(int initial_capacity) {
for (int i = 0; i < initial_capacity; i++) {
sa->fd_to_index[i] = -1;
sa->index_to_fd[i] = -1;
sa->active_indices[i] = -1;
sa->sockets[i].fd = -1;
sa->sockets[i].active = 0;
}
@ -90,6 +94,7 @@ static void socket_array_destroy(struct socket_array* sa) {
free(sa->sockets);
free(sa->fd_to_index);
free(sa->index_to_fd);
free(sa->active_indices);
free(sa);
}
@ -103,12 +108,14 @@ static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t r
struct socket_node* new_sockets = realloc(sa->sockets, new_capacity * sizeof(struct socket_node));
int* new_fd_to_index = realloc(sa->fd_to_index, new_capacity * sizeof(int));
int* new_index_to_fd = realloc(sa->index_to_fd, new_capacity * sizeof(int));
int* new_active_indices = realloc(sa->active_indices, new_capacity * sizeof(int));
if (!new_sockets || !new_fd_to_index || !new_index_to_fd) {
if (!new_sockets || !new_fd_to_index || !new_index_to_fd || !new_active_indices) {
// Allocation failed
free(new_sockets);
free(new_fd_to_index);
free(new_index_to_fd);
free(new_active_indices);
return -1;
}
@ -116,6 +123,7 @@ static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t r
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;
}
@ -123,6 +131,7 @@ static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t r
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;
}
@ -150,6 +159,7 @@ static int socket_array_add(struct socket_array* sa, int fd, socket_callback_t r
sa->fd_to_index[fd] = index;
sa->index_to_fd[index] = fd;
sa->active_indices[sa->count] = index; // Add to active list
sa->count++;
if (fd > sa->max_fd) sa->max_fd = fd;
@ -168,6 +178,17 @@ static int socket_array_remove(struct socket_array* sa, int fd) {
sa->sockets[index].fd = -1;
sa->fd_to_index[fd] = -1;
sa->index_to_fd[index] = -1;
// Remove from active_indices by swapping with last element
// Find position in active_indices
for (int i = 0; i < sa->count; i++) {
if (sa->active_indices[i] == index) {
// Swap with last element
sa->active_indices[i] = sa->active_indices[sa->count - 1];
sa->active_indices[sa->count - 1] = -1;
break;
}
}
sa->count--;
return 0;
@ -252,6 +273,7 @@ static void process_timeouts(struct UASYNC* ua) {
node->ua->timer_free_count++;
}
free(node);
break;
}
}
@ -292,7 +314,7 @@ void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_c
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);
// 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 = malloc(sizeof(struct timeout_node));
if (!node) {
@ -340,7 +362,7 @@ err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) {
// Successfully marked as deleted - free will happen lazily in heap
node->cancelled = 1;
node->callback = NULL;
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: successfully cancelled timer %p from heap", node);
// DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "uasync_cancel_timeout: successfully cancelled timer %p from heap", node);
return ERR_OK;
}
@ -360,6 +382,7 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s
if (index < 0) return NULL;
ua->socket_alloc_count++;
ua->poll_fds_dirty = 1; // Mark poll_fds as needing rebuild
// Return pointer to the socket node as ID
return &ua->sockets->sockets[index];
@ -374,16 +397,69 @@ err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) {
int ret = socket_array_remove(ua->sockets, node->fd);
if (ret == 0) {
ua->socket_free_count++;
ua->poll_fds_dirty = 1; // Mark poll_fds as needing rebuild
return ERR_OK;
}
return ERR_FAIL;
}
// Helper function to rebuild cached pollfd array
static void rebuild_poll_fds(struct UASYNC* ua) {
if (!ua || !ua->sockets) return;
int socket_count = ua->sockets->count;
int wakeup_fd_present = ua->wakeup_initialized && ua->wakeup_pipe[0] >= 0;
int total_fds = socket_count + wakeup_fd_present;
// Ensure poll_fds capacity is sufficient
if (total_fds > ua->poll_fds_capacity) {
int new_capacity = total_fds * 2;
if (new_capacity < 16) new_capacity = 16;
struct pollfd* new_poll_fds = realloc(ua->poll_fds, sizeof(struct pollfd) * new_capacity);
if (!new_poll_fds) return; // Keep old allocation on failure
ua->poll_fds = new_poll_fds;
ua->poll_fds_capacity = new_capacity;
}
int idx = 0;
// Add wakeup fd first if present
if (wakeup_fd_present) {
ua->poll_fds[idx].fd = ua->wakeup_pipe[0];
ua->poll_fds[idx].events = POLLIN;
ua->poll_fds[idx].revents = 0;
idx++;
}
// Add socket fds using active_indices for O(1) traversal
for (int i = 0; i < socket_count; i++) {
int socket_array_idx = ua->sockets->active_indices[i];
struct socket_node* cur = &ua->sockets->sockets[socket_array_idx];
ua->poll_fds[idx].fd = cur->fd;
ua->poll_fds[idx].events = 0;
ua->poll_fds[idx].revents = 0;
if (cur->read_cbk) ua->poll_fds[idx].events |= POLLIN;
if (cur->write_cbk) ua->poll_fds[idx].events |= POLLOUT;
if (cur->except_cbk) ua->poll_fds[idx].events |= POLLPRI;
idx++;
}
ua->poll_fds_count = total_fds;
ua->poll_fds_dirty = 0;
}
// Instance version
void uasync_poll(struct UASYNC* ua, int timeout_tb) {
if (!ua) return;
if (!ua->sockets || !ua->timeout_heap) return;
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "poll");
// Handle negative or zero timeout
if (timeout_tb < 0) timeout_tb = -1; // Infinite wait
@ -440,79 +516,45 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
int wakeup_fd_present = ua->wakeup_initialized && ua->wakeup_pipe[0] >= 0;
int total_fds = socket_count + wakeup_fd_present;
// Allocate pollfds and node pointers
struct pollfd* fds = malloc(sizeof(struct pollfd) * total_fds);
struct socket_node** nodes = wakeup_fd_present ? malloc(sizeof(struct socket_node*) * (total_fds - 1)) : malloc(sizeof(struct socket_node*) * total_fds);
if (!fds || !nodes) {
free(fds);
free(nodes);
return;
// Rebuild poll_fds if dirty or not allocated
if (ua->poll_fds_dirty || !ua->poll_fds) {
rebuild_poll_fds(ua);
}
int idx = 0;
// Add wakeup fd first if present
if (wakeup_fd_present) {
fds[idx].fd = ua->wakeup_pipe[0];
fds[idx].events = POLLIN;
fds[idx].revents = 0;
idx++;
}
/* Add socket fds using efficient array traversal */
int node_idx = 0;
for (int i = 0; i < ua->sockets->capacity && node_idx < socket_count; i++) {
if (ua->sockets->sockets[i].active) {
struct socket_node* cur = &ua->sockets->sockets[i];
fds[idx].fd = cur->fd;
fds[idx].events = 0;
fds[idx].revents = 0;
if (cur->read_cbk) fds[idx].events |= POLLIN;
if (cur->write_cbk) fds[idx].events |= POLLOUT;
if (cur->except_cbk) fds[idx].events |= POLLPRI;
if (nodes) {
nodes[node_idx] = cur;
}
idx++;
node_idx++;
}
// Ensure poll_fds_count matches current state (in case sockets changed without dirty flag)
if (ua->poll_fds_count != total_fds) {
rebuild_poll_fds(ua);
}
/* Call poll */
int ret = poll(fds, total_fds, timeout_ms);
/* Call poll with cached fds */
int ret = poll(ua->poll_fds, ua->poll_fds_count, timeout_ms);
if (ret < 0) {
if (errno == EINTR) {
free(fds);
free(nodes);
return;
}
perror("poll");
free(fds);
free(nodes);
return;
}
/* Process socket events first to give sockets higher priority */
if (ret > 0) {
for (int i = 0; i < total_fds; i++) {
if (fds[i].revents == 0) continue;
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 (fds[i].revents & POLLIN) {
if (ua->poll_fds[i].revents & POLLIN) {
drain_wakeup_pipe(ua);
}
continue;
}
/* Socket event */
int socket_idx = i - wakeup_fd_present;
struct socket_node* node = nodes[socket_idx];
/* Socket event - lookup by fd */
struct socket_node* node = socket_array_get(ua->sockets, ua->poll_fds[i].fd);
if (!node) continue; // Socket may have been removed
/* Check for error conditions first */
if (fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) {
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);
@ -520,21 +562,21 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
}
/* Exceptional data (out-of-band) */
if (fds[i].revents & POLLPRI) {
if (ua->poll_fds[i].revents & POLLPRI) {
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
/* Read readiness */
if (fds[i].revents & POLLIN) {
if (ua->poll_fds[i].revents & POLLIN) {
if (node->read_cbk) {
node->read_cbk(node->fd, node->user_data);
}
}
/* Write readiness */
if (fds[i].revents & POLLOUT) {
if (ua->poll_fds[i].revents & POLLOUT) {
if (node->write_cbk) {
node->write_cbk(node->fd, node->user_data);
}
@ -544,9 +586,6 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
/* Process timeouts that may have expired during poll or socket processing */
process_timeouts(ua);
free(fds);
free(nodes);
}
@ -715,6 +754,9 @@ void uasync_destroy(struct UASYNC* ua, int close_fds) {
close(ua->wakeup_pipe[1]);
}
// Free cached poll_fds
free(ua->poll_fds);
// Final leak check
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 after cleanup: timers %zu/%zu, sockets %zu/%zu",

31
lib/u_async.h

@ -19,19 +19,24 @@ typedef int err_t;
#define ERR_OK 0
#define ERR_FAIL -1
// Uasync instance structure
struct UASYNC {
TimeoutHeap* timeout_heap; // Heap for timeout management
struct socket_array* sockets; // Array-based socket management
// Debug counters for memory allocation tracking
size_t timer_alloc_count;
size_t timer_free_count;
size_t socket_alloc_count;
size_t socket_free_count;
// Wakeup pipe for interrupting poll
int wakeup_pipe[2]; // [0] read, [1] write
int wakeup_initialized;
};
// Uasync instance structure
struct UASYNC {
TimeoutHeap* timeout_heap; // Heap for timeout management
struct socket_array* sockets; // Array-based socket management
// Debug counters for memory allocation tracking
size_t timer_alloc_count;
size_t timer_free_count;
size_t socket_alloc_count;
size_t socket_free_count;
// Wakeup pipe for interrupting poll
int wakeup_pipe[2]; // [0] read, [1] write
int wakeup_initialized;
// Cached pollfd array for optimized poll operations
struct pollfd* poll_fds; // Pre-allocated pollfd array
int poll_fds_capacity; // Current capacity of poll_fds
int poll_fds_count; // Number of valid entries in poll_fds
int poll_fds_dirty; // 1 if poll_fds needs rebuild
};
// Type definitions
typedef struct UASYNC uasync_t;

Loading…
Cancel
Save