Browse Source

fix uasync: posted_tasks delivery on Linux, create cleanup leaks, atomic array expand

congestion
Evgeny 5 months ago
parent
commit
4d5afe5ae9
  1. 144
      lib/u_async.c

144
lib/u_async.c

@ -129,39 +129,46 @@ 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) {
// 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;
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;
}
// Check if FD already exists
@ -839,7 +846,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
for (int i = 0; i < n_events; i++) {
if (events[i].data.fd == -1) {
if (events[i].events & EPOLLIN) {
drain_wakeup_pipe(ua);
handle_wakeup(ua);
}
continue;
}
@ -1164,7 +1171,7 @@ 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) {
drain_wakeup_pipe(ua);
handle_wakeup(ua);
}
continue;
}
@ -1376,27 +1383,27 @@ struct UASYNC* uasync_create(void) {
#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");
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;
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;
}
ua->wakeup_pipe[0] = (int)(intptr_t)r;
@ -1413,10 +1420,9 @@ struct UASYNC* uasync_create(void) {
#else
// POSIX pipe
if (pipe(ua->wakeup_pipe) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno));
u_free(ua);
return NULL;
if (pipe(ua->wakeup_pipe) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_UASYNC, "pipe() failed: %s", strerror(errno));
goto create_cleanup;
}
ua->wakeup_initialized = 1;
@ -1428,10 +1434,20 @@ struct UASYNC* uasync_create(void) {
}
pthread_mutex_init(&ua->posted_lock, NULL);
#endif
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Creating TH3...");
return ua;
#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;
}
// Print all resources for debugging

Loading…
Cancel
Save