Browse Source

fix stale pointer after realloc in uasync socket/event handlers

congestion
Evgeny 5 months ago
parent
commit
1a4b0a9de6
  1. 484
      lib/u_async.c

484
lib/u_async.c

@ -661,28 +661,25 @@ void* uasync_add_socket(struct UASYNC* ua, int fd, socket_callback_t read_cbk, s
if (read_cbk) ev.events |= EPOLLIN;
if (write_cbk) ev.events |= EPOLLOUT;
// Use level-triggered mode (default) for compatibility with UDP sockets
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];
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;
}
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;
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;
#if HAS_EPOLL
// Remove from epoll if using epoll
@ -709,32 +706,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;
#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.ptr = &ua->sockets->sockets[index];
// 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
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];
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;
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;
}
// Remove socket by socket_t
@ -836,56 +832,73 @@ 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);
}
}
/* Read readiness - use appropriate callback based on socket type */
if (events[i].events & EPOLLIN) {
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 (events[i].events & EPOLLOUT) {
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);
}
}
}
}
}
// 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) {
if (events[i].events & EPOLLIN) {
drain_wakeup_pipe(ua);
}
continue;
}
int fd = events[i].data.fd;
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 (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 (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);
}
}
}
}
}
}
#endif
// Instance version
@ -1039,55 +1052,87 @@ 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;
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);
}
}
}
}
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);
}
}
}
}
}
}
#else
// On non-Windows, use poll()
@ -1111,68 +1156,91 @@ 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) {
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);
}
}
}
}
/* 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;
}
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);
}
}
}
}
}
}
#endif
@ -1294,7 +1362,7 @@ struct UASYNC* uasync_create(void) {
if (ua->wakeup_initialized) {
struct epoll_event ev;
ev.events = EPOLLIN;
ev.data.ptr = NULL; // NULL ptr indicates wakeup fd
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));
}

Loading…
Cancel
Save