|
|
|
|
@ -11,6 +11,14 @@
|
|
|
|
|
#include <limits.h> |
|
|
|
|
#include <fcntl.h> |
|
|
|
|
|
|
|
|
|
// Platform-specific includes
|
|
|
|
|
#ifdef __linux__ |
|
|
|
|
#include <sys/epoll.h> |
|
|
|
|
#define HAS_EPOLL 1 |
|
|
|
|
#else |
|
|
|
|
#define HAS_EPOLL 0 |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Timeout node with safe cancellation
|
|
|
|
|
@ -384,6 +392,26 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s
|
|
|
|
|
ua->socket_alloc_count++; |
|
|
|
|
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; |
|
|
|
|
// Use edge-triggered mode for better performance
|
|
|
|
|
ev.events |= EPOLLET; |
|
|
|
|
ev.data.ptr = &ua->sockets->sockets[index]; |
|
|
|
|
|
|
|
|
|
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; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
// Return pointer to the socket node as ID
|
|
|
|
|
return &ua->sockets->sockets[index]; |
|
|
|
|
} |
|
|
|
|
@ -394,7 +422,16 @@ err_t uasync_remove_socket(struct UASYNC* ua, void* s_id) {
|
|
|
|
|
struct socket_node* node = (struct socket_node*)s_id; |
|
|
|
|
if (!node->active || node->fd < 0) return ERR_FAIL; |
|
|
|
|
|
|
|
|
|
int ret = socket_array_remove(ua->sockets, node->fd); |
|
|
|
|
int fd = node->fd; |
|
|
|
|
|
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
// Remove from epoll if using epoll
|
|
|
|
|
if (ua->use_epoll && ua->epoll_fd >= 0) { |
|
|
|
|
epoll_ctl(ua->epoll_fd, EPOLL_CTL_DEL, fd, NULL); |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
int ret = socket_array_remove(ua->sockets, fd); |
|
|
|
|
if (ret == 0) { |
|
|
|
|
ua->socket_free_count++; |
|
|
|
|
ua->poll_fds_dirty = 1; // Mark poll_fds as needing rebuild
|
|
|
|
|
@ -454,6 +491,48 @@ static void rebuild_poll_fds(struct UASYNC* ua) {
|
|
|
|
|
ua->poll_fds_dirty = 0; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// 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.ptr is NULL)
|
|
|
|
|
if (events[i].data.ptr == NULL) { |
|
|
|
|
if (events[i].events & EPOLLIN) { |
|
|
|
|
drain_wakeup_pipe(ua); |
|
|
|
|
} |
|
|
|
|
continue; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// Socket event
|
|
|
|
|
struct socket_node* node = (struct socket_node*)events[i].data.ptr; |
|
|
|
|
if (!node || !node->active) continue; |
|
|
|
|
|
|
|
|
|
/* Check for error conditions first */ |
|
|
|
|
if (events[i].events & (EPOLLERR | EPOLLHUP)) { |
|
|
|
|
if (node->except_cbk) { |
|
|
|
|
node->except_cbk(node->fd, node->user_data); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Exceptional data (out-of-band) - epoll doesn't have POLLPRI equivalent */ |
|
|
|
|
|
|
|
|
|
/* Read readiness */ |
|
|
|
|
if (events[i].events & EPOLLIN) { |
|
|
|
|
if (node->read_cbk) { |
|
|
|
|
node->read_cbk(node->fd, node->user_data); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Write readiness */ |
|
|
|
|
if (events[i].events & EPOLLOUT) { |
|
|
|
|
if (node->write_cbk) { |
|
|
|
|
node->write_cbk(node->fd, node->user_data); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
// Instance version
|
|
|
|
|
void uasync_poll(struct UASYNC* ua, int timeout_tb) { |
|
|
|
|
if (!ua) return; |
|
|
|
|
@ -512,6 +591,34 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
|
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
// Use epoll on Linux if available
|
|
|
|
|
if (ua->use_epoll && ua->epoll_fd >= 0) { |
|
|
|
|
struct epoll_event events[64]; // Stack-allocated array for events
|
|
|
|
|
int max_events = 64; |
|
|
|
|
|
|
|
|
|
int ret = epoll_wait(ua->epoll_fd, events, max_events, timeout_ms); |
|
|
|
|
if (ret < 0) { |
|
|
|
|
if (errno == EINTR) { |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
perror("epoll_wait"); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Process socket events */ |
|
|
|
|
if (ret > 0) { |
|
|
|
|
process_epoll_events(ua, events, ret); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Process timeouts that may have expired during poll or socket processing */ |
|
|
|
|
process_timeouts(ua); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
// Fallback to poll() for non-Linux or if epoll failed
|
|
|
|
|
|
|
|
|
|
// Include wakeup pipe if initialized
|
|
|
|
|
int wakeup_fd_present = ua->wakeup_initialized && ua->wakeup_pipe[0] >= 0; |
|
|
|
|
int total_fds = socket_count + wakeup_fd_present; |
|
|
|
|
@ -641,6 +748,28 @@ struct UASYNC* uasync_create(void) {
|
|
|
|
|
// Set callback to free timeout nodes and update counters
|
|
|
|
|
timeout_heap_set_free_callback(ua->timeout_heap, ua, timeout_node_free_callback); |
|
|
|
|
|
|
|
|
|
// Initialize epoll on Linux
|
|
|
|
|
ua->epoll_fd = -1; |
|
|
|
|
ua->use_epoll = 0; |
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
ua->epoll_fd = epoll_create1(EPOLL_CLOEXEC); |
|
|
|
|
if (ua->epoll_fd >= 0) { |
|
|
|
|
ua->use_epoll = 1; |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_UASYNC, "Using epoll for socket monitoring"); |
|
|
|
|
// Add wakeup pipe to epoll
|
|
|
|
|
if (ua->wakeup_initialized) { |
|
|
|
|
struct epoll_event ev; |
|
|
|
|
ev.events = EPOLLIN; |
|
|
|
|
ev.data.ptr = NULL; // NULL ptr indicates wakeup fd
|
|
|
|
|
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)); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} else { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_UASYNC, "Failed to create epoll fd, falling back to poll: %s", strerror(errno)); |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
return ua; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -757,6 +886,13 @@ void uasync_destroy(struct UASYNC* ua, int close_fds) {
|
|
|
|
|
// Free cached poll_fds
|
|
|
|
|
free(ua->poll_fds); |
|
|
|
|
|
|
|
|
|
// Close epoll fd on Linux
|
|
|
|
|
#if HAS_EPOLL |
|
|
|
|
if (ua->epoll_fd >= 0) { |
|
|
|
|
close(ua->epoll_fd); |
|
|
|
|
} |
|
|
|
|
#endif |
|
|
|
|
|
|
|
|
|
// 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", |
|
|
|
|
|