@ -4,7 +4,8 @@
# include "platform_compat.h"
# include "debug_config.h"
# include "mem.h"
# include "memory_pool.h"
# include "memory_pool.h"
# include "twheel.h"
# include <stdio.h>
# include <string.h>
# include <stdlib.h>
@ -32,13 +33,12 @@
// Timeout node with safe cancellation
struct timeout_node {
struct twheel to ; // MUST be first — embedded timing wheel entry
char name [ 16 ] ;
void * arg ;
timeout_callback_t callback ;
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
@ -287,19 +287,52 @@ static struct socket_node* socket_array_get_by_sock(struct socket_array* sa, soc
if ( index = = - 1 | | ! sa - > sockets [ index ] . active ) return NULL ;
if ( sa - > sockets [ index ] . type ! = SOCKET_NODE_TYPE_SOCK ) return NULL ;
return & sa - > sockets [ index ] ;
}
// Callback to u_free timeout node and update counters
static void timeout_node_free_callback ( void * user_data , void * data ) {
struct UASYNC * ua = ( struct UASYNC * ) user_data ;
struct timeout_node * node = ( struct timeout_node * ) data ;
( void ) node ; // Not used directly, but keep for consistency
ua - > timer_free_count + + ;
memory_pool_free ( ua - > timeout_pool , data ) ;
}
// Helper to get current time
return & sa - > sockets [ index ] ;
}
// ------ socket_node dispatch helpers (eliminate FD/SOCK type duplication) ------
static inline int socket_node_fd ( const struct socket_node * n ) {
if ( n - > type = = SOCKET_NODE_TYPE_SOCK ) {
# ifdef _WIN32
return ( int ) ( intptr_t ) n - > sock ;
# else
return n - > sock ;
# endif
}
return n - > fd ;
}
static inline void socket_node_dispatch_read ( const struct socket_node * n ) {
if ( n - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( n - > read_cbk_sock ) n - > read_cbk_sock ( n - > sock , n - > user_data ) ;
} else {
if ( n - > read_cbk ) n - > read_cbk ( n - > fd , n - > user_data ) ;
}
}
static inline void socket_node_dispatch_write ( const struct socket_node * n ) {
if ( n - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( n - > write_cbk_sock ) n - > write_cbk_sock ( n - > sock , n - > user_data ) ;
} else {
if ( n - > write_cbk ) n - > write_cbk ( n - > fd , n - > user_data ) ;
}
}
static inline short socket_node_events ( const struct socket_node * n ) {
short events = 0 ;
if ( n - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( n - > read_cbk_sock & & n - > enable_read ) events | = POLLIN ;
if ( n - > write_cbk_sock & & n - > enable_write ) events | = POLLOUT ;
} else {
if ( n - > read_cbk & & n - > enable_read ) events | = POLLIN ;
if ( n - > write_cbk & & n - > enable_write ) events | = POLLOUT ;
}
if ( n - > except_cbk ) events | = POLLPRI ;
return events ;
}
// Simplified timeout handling without reference counting
static void get_current_time ( struct timeval * tv ) {
# ifdef _WIN32
// Для Windows используем GetTickCount или другие механизмы
@ -358,17 +391,6 @@ uint64_t get_time_us(void) {
// Drain wakeup pipe - read all available bytes
static void drain_wakeup_pipe ( struct UASYNC * ua ) {
if ( ! ua | | ! ua - > wakeup_initialized ) return ;
char buf [ 64 ] ;
while ( 1 ) {
ssize_t n = read ( ua - > wakeup_pipe [ 0 ] , buf , sizeof ( buf ) ) ;
if ( n < = 0 ) break ;
}
}
// Process posted tasks (lock-u_free during execution)
static void process_posted_tasks ( struct UASYNC * ua ) {
if ( ! ua ) return ;
@ -390,7 +412,7 @@ static void process_posted_tasks(struct UASYNC* ua) {
# endif
while ( list ) {
DEBUG_DEBUG ( DEBUG_CATEGORY_T UN , " POSTed task get " ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_UASY NC , " POSTed task get " ) ;
struct posted_task * t = list ;
list = list - > next ;
@ -435,8 +457,8 @@ static uint64_t timeval_to_ms(const struct timeval* tv) {
// Simplified timeout handling without reference counting
// Simplified timeout handling without reference counting
// Process expired timeouts with safe cancellation
static void process_timeouts ( struct UASYNC * ua ) {
if ( ! ua ) return ;
@ -460,63 +482,40 @@ static void process_timeouts(struct UASYNC* ua) {
memory_pool_free ( ua - > timeout_pool , node ) ;
}
if ( ! ua - > timeout_heap ) return ;
if ( ! ua - > twheel ) return ;
struct timeval now_tv ;
get_current_time ( & now_tv ) ;
uint64_t now_ms = timeval_to_ms ( & now_tv ) ;
while ( 1 ) {
TimeoutEntry entry ;
if ( timeout_heap_peek ( ua - > timeout_heap , & entry ) ! = 0 ) break ;
if ( entry . expiration > now_ms ) break ;
// Pop the expired timeout
timeout_heap_pop ( ua - > timeout_heap , & entry ) ;
struct timeout_node * node = ( struct timeout_node * ) entry . data ;
// Update timing wheel to current time — moves expired to expired queue
twheels_update ( ua - > twheel , now_ms ) ;
// Drain all expired timeouts
struct twheel * to ;
while ( ( to = twheels_get ( ua - > twheel ) ) ) {
struct timeout_node * node = ( struct timeout_node * ) to ;
if ( node & & node - > callback ) {
DEBUG_DEBUG ( DEBUG_CATEGORY_UASYNC , " timer→%s expired " , node - > name [ 0 ] ? node - > name : " " ) ;
node - > callback ( node - > arg ) ;
}
// Always u_free the node after processing
if ( node & & node - > ua ) {
node - > ua - > timer_free_count + + ;
}
memory_pool_free ( ua - > timeout_pool , node ) ;
continue ; // Process next expired timeout
}
}
// Compute time to next timeout
static void get_next_timeout ( struct UASYNC * ua , struct timeval * tv ) {
if ( ! ua | | ! ua - > timeout_heap ) {
tv - > tv_sec = 0 ;
tv - > tv_usec = 0 ;
return ;
}
TimeoutEntry entry ;
if ( timeout_heap_peek ( ua - > timeout_heap , & entry ) ! = 0 ) {
tv - > tv_sec = 0 ;
tv - > tv_usec = 0 ;
return ;
}
struct timeval now_tv ;
get_current_time ( & now_tv ) ;
uint64_t now_ms = timeval_to_ms ( & now_tv ) ;
if ( entry . expiration < = now_ms ) {
tv - > tv_sec = 0 ;
tv - > tv_usec = 0 ;
return ;
}
uint64_t delta_ms = entry . expiration - now_ms ;
tv - > tv_sec = delta_ms / 1000 ;
tv - > tv_usec = ( delta_ms % 1000 ) * 1000 ;
// Compute time to next timeout in milliseconds
// Returns: 0 = timer already expired/no timers, ~0 = no pending timers, >0 = ms to wait
static uint64_t get_next_timeout_ms ( struct UASYNC * ua ) {
if ( ! ua | | ! ua - > twheel ) return 0 ;
twheel_t t = twheels_timeout ( ua - > twheel ) ;
if ( t = = 0 ) return 0 ; // expired timer exists
if ( t = = ~ TWHEEL_C ( 0 ) ) return ~ TWHEEL_C ( 0 ) ; // no pending timers
return ( uint64_t ) t ;
}
@ -524,9 +523,7 @@ static void get_next_timeout(struct UASYNC* ua, struct timeval* tv) {
// Instance version
void * uasync_set_timeout ( struct UASYNC * ua , int timeout_tb , void * arg , timeout_callback_t callback , const char * name ) {
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);
if ( ! ua - > twheel ) return NULL ;
struct timeout_node * node = memory_pool_alloc ( ua - > timeout_pool ) ;
if ( ! node ) {
@ -544,24 +541,19 @@ void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_c
node - > arg = arg ;
node - > callback = callback ;
node - > ua = ua ;
node - > heap_index = SIZE_MAX ;
// Calculate expiration time in milliseconds
// Calculate expiration time in milliseconds (absolute monotonic)
struct timeval now ;
get_current_time ( & now ) ;
timeval_add_tb ( & now , timeout_tb ) ;
node - > expiration_ms = timeval_to_ms ( & now ) ;
uint64_t expiration_ms = timeval_to_ms ( & now ) ;
// Initialize wheel entry and add to timing wheel
twheel_init ( & node - > to , TWHEEL_ABS ) ;
twheels_add ( ua - > twheel , & node - > to , expiration_ms ) ;
// Add to heap
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
return NULL ;
}
return node ;
}
}
// Immediate execution in next mainloop (FIFO order)
void * uasync_call_soon ( struct UASYNC * ua , void * user_arg , timeout_callback_t callback ) {
@ -579,9 +571,8 @@ void* uasync_call_soon(struct UASYNC* ua, void* user_arg, timeout_callback_t cal
node - > arg = user_arg ;
node - > callback = callback ;
node - > ua = ua ;
node - > expiration_ms = 0 ;
node - > next = NULL ;
node - > heap_index = SIZE_MAX ;
memset ( & node - > to , 0 , sizeof ( node - > to ) ) ;
// FIFO: добавляем в конец очереди
if ( ua - > immediate_queue_tail ) {
@ -609,30 +600,26 @@ err_t uasync_call_soon_cancel(struct UASYNC* ua, void* t_id) {
// Instance version
// Instance version
err_t uasync_cancel_timeout ( struct UASYNC * ua , void * t_id ) {
if ( ! ua | | ! t_id | | ! ua - > timeout_heap ) {
DEBUG_ERROR ( DEBUG_CATEGORY_TIMERS , " uasync_cancel_timeout: invalid parameters ua=%p, t_id=%p, heap =%p " ,
ua , t_id , ua ? ua - > timeout_heap : NULL ) ;
if ( ! ua | | ! t_id | | ! ua - > twheel ) {
DEBUG_ERROR ( DEBUG_CATEGORY_TIMERS , " uasync_cancel_timeout: invalid parameters ua=%p, t_id=%p, twheel =%p " ,
ua , t_id , ua ? ua - > twheel : NULL ) ;
return ERR_FAIL ;
}
struct timeout_node * node = ( struct timeout_node * ) t_id ;
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 ) ;
if ( ! twheel_pending ( & node - > to ) ) {
DEBUG_DEBUG ( DEBUG_CATEGORY_TIMERS , " uasync_cancel_timeout: not pending: ua=%p, t_id=%p " , ua , t_id ) ;
return ERR_FAIL ;
}
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 ;
twheels_del ( ua - > twheel , & node - > to ) ;
node - > callback = NULL ;
ua - > timer_free_count + + ;
memory_pool_free ( ua - > timeout_pool , node ) ;
return ERR_OK ;
}
@ -765,36 +752,27 @@ err_t uasync_remove_socket_t(struct UASYNC* ua, socket_t sock) {
return ERR_FAIL ;
}
err_t uasync_set_socket_read ( struct UASYNC * ua , void * s_id , int enable ) {
static err_t socket_set_event ( struct UASYNC * ua , void * s_id , int is_rea d , int enable ) {
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 val = enable ? 1 : 0 ;
if ( node - > enable_read = = val ) return ERR_OK ;
node - > enable_read = val ;
int * field = is_read ? & node - > enable_read : & node - > enable_write ;
if ( * field = = val ) return ERR_OK ;
* field = val ;
# if HAS_EPOLL
if ( ua - > use_epoll & & ua - > epoll_fd > = 0 ) {
struct epoll_event ev ;
ev . events = 0 ;
if ( node - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( node - > read_cbk_sock & & node - > enable_read ) ev . events | = EPOLLIN ;
if ( node - > write_cbk_sock & & node - > enable_write ) ev . events | = EPOLLOUT ;
ev . data . fd = node - > sock ;
} else {
if ( node - > read_cbk & & node - > enable_read ) ev . events | = EPOLLIN ;
if ( node - > write_cbk & & node - > enable_write ) ev . events | = EPOLLOUT ;
ev . data . fd = node - > fd ;
}
if ( node - > except_cbk ) ev . events | = EPOLLPRI ;
ev . events = socket_node_events ( node ) ;
ev . data . fd = socket_node_fd ( node ) ;
# ifdef _WIN32
int efd = ( int ) ( intptr_t ) ev . data . fd ;
epoll_ctl ( ua - > epoll_fd , EPOLL_CTL_MOD , ( int ) ( intptr_t ) ev . data . fd , & ev ) ;
# else
int efd = ev . data . fd ;
epoll_ctl ( ua - > epoll_fd , EPOLL_CTL_MOD , ev . data . fd , & ev ) ;
# endif
epoll_ctl ( ua - > epoll_fd , EPOLL_CTL_MOD , efd , & ev ) ;
}
# endif
@ -802,41 +780,12 @@ err_t uasync_set_socket_read(struct UASYNC* ua, void* s_id, int enable) {
return ERR_OK ;
}
err_t uasync_set_socket_write ( struct UASYNC * ua , void * s_id , int enable ) {
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 val = enable ? 1 : 0 ;
if ( node - > enable_write = = val ) return ERR_OK ;
node - > enable_write = val ;
# if HAS_EPOLL
if ( ua - > use_epoll & & ua - > epoll_fd > = 0 ) {
struct epoll_event ev ;
ev . events = 0 ;
if ( node - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( node - > read_cbk_sock & & node - > enable_read ) ev . events | = EPOLLIN ;
if ( node - > write_cbk_sock & & node - > enable_write ) ev . events | = EPOLLOUT ;
ev . data . fd = node - > sock ;
} else {
if ( node - > read_cbk & & node - > enable_read ) ev . events | = EPOLLIN ;
if ( node - > write_cbk & & node - > enable_write ) ev . events | = EPOLLOUT ;
ev . data . fd = node - > fd ;
}
if ( node - > except_cbk ) ev . events | = EPOLLPRI ;
# ifdef _WIN32
int efd = ( int ) ( intptr_t ) ev . data . fd ;
# else
int efd = ev . data . fd ;
# endif
epoll_ctl ( ua - > epoll_fd , EPOLL_CTL_MOD , efd , & ev ) ;
}
# endif
err_t uasync_set_socket_read ( struct UASYNC * ua , void * s_id , int enable ) {
return socket_set_event ( ua , s_id , 1 , enable ) ;
}
ua - > poll_fds_dirty = 1 ;
return ERR_OK ;
err_t uasync_set_socket_write ( struct UASYNC * ua , void * s_id , int enable ) {
return socket_set_event ( ua , s_id , 0 , enable ) ;
}
// Helper function to rebuild cached pollfd array
@ -869,37 +818,16 @@ static void rebuild_poll_fds(struct UASYNC* ua) {
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 ] ;
// Handle socket_t vs int fd
if ( cur - > type = = SOCKET_NODE_TYPE_SOCK ) {
// socket_t - cast to int for pollfd
# ifdef _WIN32
ua - > poll_fds [ idx ] . fd = ( int ) ( intptr_t ) cur - > sock ;
# else
ua - > poll_fds [ idx ] . fd = cur - > sock ;
# endif
} else {
// Regular fd
ua - > poll_fds [ idx ] . fd = cur - > fd ;
}
ua - > poll_fds [ idx ] . events = 0 ;
ua - > poll_fds [ idx ] . revents = 0 ;
if ( cur - > type = = SOCKET_NODE_TYPE_SOCK ) {
if ( cur - > read_cbk_sock & & cur - > enable_read ) ua - > poll_fds [ idx ] . events | = POLLIN ;
if ( cur - > write_cbk_sock & & cur - > enable_write ) ua - > poll_fds [ idx ] . events | = POLLOUT ;
} else {
if ( cur - > read_cbk & & cur - > enable_read ) ua - > poll_fds [ idx ] . events | = POLLIN ;
if ( cur - > write_cbk & & cur - > enable_write ) ua - > poll_fds [ idx ] . events | = POLLOUT ;
}
if ( cur - > write_cbk & & cur - > enable_write ) ua - > poll_fds [ idx ] . events | = POLLOUT ;
if ( cur - > except_cbk ) ua - > poll_fds [ idx ] . events | = POLLPRI ;
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 = socket_node_fd ( cur ) ;
ua - > poll_fds [ idx ] . events = socket_node_events ( cur ) ;
ua - > poll_fds [ idx ] . revents = 0 ;
idx + + ;
}
ua - > poll_fds_count = total_fds ;
@ -913,7 +841,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
// Check if this is the wakeup fd (data.fd is -1)
if ( events [ i ] . data . fd < 0 ) {
if ( events [ i ] . events & EPOLLIN ) {
drain_wakeup_pipe ( ua ) ;
handle_wakeup ( ua ) ;
}
continue ;
}
@ -936,99 +864,59 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
/* Check for error conditions first */
if ( events [ i ] . events & ( EPOLLERR | EPOLLHUP ) ) {
if ( local_except ) {
local_except ( local_fd , local_ud ) ;
}
if ( local_except ) local_except ( local_fd , local_ud ) ;
}
/* Read readiness - use appropriate callback based on socket type */
/* Read readiness */
if ( events [ i ] . events & EPOLLIN ) {
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 ) ;
}
}
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 */
/* Write readiness */
if ( events [ i ] . events & EPOLLOUT ) {
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 ) ;
}
}
if ( local_write_sock ) local_write_sock ( local_sock , local_ud ) ;
} else if ( local_write ) local_write ( local_fd , local_ud ) ;
}
}
}
# endif
// 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_UASYNC , " poll(%d sockets, %zu timers, timeout=%d.%dms) " ,
ua - > sockets - > count , ua - > timeout_heap - > size , timeout_tb > = 0 ? timeout_tb / 10000 : - 1 ,
timeout_tb > = 0 ? ( timeout_tb % 10000 ) / 10 : 0 ) ;
// Handle negative or zero timeout
if ( timeout_tb < 0 ) timeout_tb = - 1 ; // Infinite wait
else if ( timeout_tb = = 0 ) timeout_tb = 0 ; // No wait
// Get next timeout
struct timeval next_timeout ;
get_next_timeout ( ua , & next_timeout ) ;
// Convert requested timeout to timeval
struct timeval req_timeout = { 0 } ;
if ( timeout_tb > = 0 ) {
req_timeout . tv_sec = timeout_tb / 10000 ;
req_timeout . tv_usec = ( timeout_tb % 10000 ) * 100 ;
}
struct timeval poll_timeout ;
// Instance version
void uasync_poll ( struct UASYNC * ua , int timeout_tb ) {
if ( ! ua ) return ;
if ( ! ua - > sockets | | ! ua - > twheel ) return ;
DEBUG_DEBUG ( DEBUG_CATEGORY_UASYNC , " poll(%d sockets, timeout=%d.%dms) " ,
ua - > sockets - > count , timeout_tb > = 0 ? timeout_tb / 10000 : - 1 ,
timeout_tb > = 0 ? ( timeout_tb % 10000 ) / 10 : 0 ) ;
uint64_t next_ms = get_next_timeout_ms ( ua ) ; // 0=expired, ~0=no timers, >0=ms to wait
int timeout_ms ;
if ( timeout_tb < 0 ) {
poll_timeout = next_timeout ;
// Infinite wait or until next timer
if ( next_ms = = 0 ) timeout_ms = 0 ; // Expired timer, don't wait
else if ( next_ms ! = ~ TWHEEL_C ( 0 ) ) timeout_ms = ( int ) next_ms ; // Wait for next timer
else timeout_ms = - 1 ; // No timers, infinite wait
} else {
if ( next_timeout . tv_sec < req_timeout . tv_sec | |
( next_timeout . tv_sec = = req_timeout . tv_sec & & next_timeout . tv_usec < req_timeout . tv_usec ) ) {
poll_timeout = next_timeout ;
} else {
poll_timeout = req_timeout ;
}
int req_ms = timeout_tb / 10 ; // timebase to ms
if ( next_ms = = 0 ) timeout_ms = 0 ; // Expired timer, don't wait
else if ( next_ms ! = ~ TWHEEL_C ( 0 ) ) timeout_ms = ( int ) ( next_ms < ( uint64_t ) req_ms ? next_ms : ( uint64_t ) req_ms ) ;
else timeout_ms = req_ms ;
}
if ( poll_timeout . tv_sec = = 0 & & poll_timeout . tv_usec = = 0 & & timeout_tb > 0 ) poll_timeout = req_timeout ;
int timeout_ms ;
if ( timeout_tb < 0 & & ( next_timeout . tv_sec > 0 | | next_timeout . tv_usec > 0 ) ) {
timeout_ms = ( poll_timeout . tv_sec * 1000 ) + ( poll_timeout . tv_usec / 1000 ) ;
} else if ( timeout_tb < 0 ) {
timeout_ms = - 1 ; // Infinite
} else {
timeout_ms = ( poll_timeout . tv_sec * 1000 ) + ( poll_timeout . tv_usec / 1000 ) ;
}
// Count active sockets
int socket_count = ua - > sockets - > count ;
if ( socket_count = = 0 & & timeout_ms = = - 1 ) {
// No sockets and infinite wait - but we have timers? Wait for timer
if ( ua - > timeout_heap - > size > 0 ) {
timeout_ms = ( poll_timeout . tv_sec * 1000 ) + ( poll_timeout . tv_usec / 1000 ) ;
} else {
// Nothing to do - return immediately
return ;
}
}
int socket_count = ua - > sockets - > count ;
if ( socket_count = = 0 & & timeout_ms = = - 1 ) {
// No sockets and infinite wait — but we might have timers
if ( twheels_pending ( ua - > twheel ) | | twheels_expired ( ua - > twheel ) )
timeout_ms = 0 ; // Let process_timeouts handle them
else
return ; // Nothing to do
}
# if HAS_EPOLL
// Use epoll on Linux if available
if ( ua - > use_epoll & & ua - > epoll_fd > = 0 ) {
@ -1143,35 +1031,12 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
if ( ! has_read & & ! has_write & & ! has_except ) continue ;
DEBUG_DEBUG ( DEBUG_CATEGORY_UASYNC , " select→fd=%d r=%d w=%d e=%d " , ( int ) s , has_read , has_write , has_except ) ;
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 ( has_except ) {
if ( node - > except_cbk ) node - > except_cbk ( node - > fd , node - > user_data ) ;
}
if ( has_read ) socket_node_dispatch_read ( node ) ;
if ( has_write ) socket_node_dispatch_write ( node ) ;
}
}
# else
@ -1202,14 +1067,6 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
if ( ua - > poll_fds [ i ] . revents = = 0 ) continue ;
DEBUG_DEBUG ( DEBUG_CATEGORY_UASYNC , " poll→fd=%d rev=0x%x " , ua - > poll_fds [ i ] . fd , ua - > poll_fds [ i ] . revents ) ;
/* 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)
@ -1218,46 +1075,18 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
}
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 ) ;
}
}
}
/* Check for error conditions first */
if ( ua - > poll_fds [ i ] . revents & ( POLLERR | POLLHUP | POLLNVAL ) ) {
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 ) ;
}
if ( ua - > poll_fds [ i ] . revents & POLLIN ) socket_node_dispatch_read ( node ) ;
if ( ua - > poll_fds [ i ] . revents & POLLOUT ) socket_node_dispatch_write ( node ) ;
}
}
# endif
@ -1267,9 +1096,6 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
}
// Put this near the top of u_async.c, after includes and before uasync_create
# ifdef _WIN32
static void wakeup_read_callback_win ( socket_t sock , void * arg ) {
( void ) sock ; // не нужен
@ -1330,8 +1156,8 @@ struct UASYNC* uasync_create(void) {
return NULL ;
}
ua - > timeout_heap = timeout_heap_create ( 16 ) ;
if ( ! ua - > timeout_heap ) {
ua - > twheel = twheels_open ( TWHEEL_mHZ ) ;
if ( ! ua - > twheel ) {
socket_array_destroy ( ua - > sockets ) ;
if ( ua - > wakeup_initialized ) {
# ifdef _WIN32
@ -1347,9 +1173,9 @@ struct UASYNC* uasync_create(void) {
}
// Initialize timeout pool
ua - > timeout_pool = memory_pool_init ( sizeof ( struct timeout_node ) ) ;
ua - > timeout_pool = memory_pool_init ( sizeof ( struct timeout_node ) , " timeout_pool " ) ;
if ( ! ua - > timeout_pool ) {
timeout_heap_destroy ( ua - > timeout_heap ) ;
twheels_close ( ua - > twheel ) ;
socket_array_destroy ( ua - > sockets ) ;
if ( ua - > wakeup_initialized ) {
# ifdef _WIN32
@ -1364,10 +1190,7 @@ struct UASYNC* uasync_create(void) {
return NULL ;
}
// Set callback to u_free timeout nodes and update counters
timeout_heap_set_free_callback ( ua - > timeout_heap , ua , timeout_node_free_callback ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " Creating TH1... " ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " Creating TH1... " ) ;
// Initialize epoll on Linux
ua - > epoll_fd = - 1 ;
ua - > use_epoll = 0 ;
@ -1479,18 +1302,11 @@ void uasync_print_resources(struct UASYNC* ua, const char* prefix) {
( ssize_t ) ( ua - > socket_alloc_count - ua - > socket_free_count ) ) ;
// Показать активные таймеры
if ( ua - > timeout_heap ) {
size_t active_timers = 0 ;
// Безопасное чтение без извлечения - просто итерируем по массиву
for ( size_t i = 0 ; i < ua - > timeout_heap - > size ; i + + ) {
if ( ! ua - > timeout_heap - > heap [ i ] . deleted ) {
active_timers + + ;
struct timeout_node * node = ( struct timeout_node * ) ua - > timeout_heap - > heap [ i ] . data ;
DEBUG_INFO ( DEBUG_CATEGORY_UASYNC , " Timer: node=%p, expires=%llu ms " ,
node , ( unsigned long long ) ua - > timeout_heap - > heap [ i ] . expiration ) ;
}
}
DEBUG_INFO ( DEBUG_CATEGORY_UASYNC , " Active timers in heap: %zu " , active_timers ) ;
if ( ua - > twheel ) {
bool has_pending = twheels_pending ( ua - > twheel ) ;
bool has_expired = twheels_expired ( ua - > twheel ) ;
DEBUG_INFO ( DEBUG_CATEGORY_UASYNC , " Timing wheel: pending=%d expired=%d curtime=%llu " ,
has_pending , has_expired , ( unsigned long long ) ua - > twheel - > curtime ) ;
}
// Показать активные сокеты
@ -1549,21 +1365,30 @@ void uasync_destroy(struct UASYNC* ua, int close_fds) {
}
ua - > immediate_queue_tail = NULL ;
// Очистить heap
if ( ua - > timeout_heap ) {
while ( 1 ) {
TimeoutEntry entry ;
if ( timeout_heap_pop ( ua - > timeout_heap , & entry ) ! = 0 ) break ;
struct timeout_node * node = ( struct timeout_node * ) entry . data ;
// Free all timer nodes (avoid double-u_free bug)
if ( node ) {
ua - > timer_free_count + + ;
memory_pool_free ( ua - > timeout_pool , node ) ;
// Очистить timing wheel
if ( ua - > twheel ) {
for ( int w = 0 ; w < 4 ; w + + ) {
for ( int s = 0 ; s < 64 ; s + + ) {
struct twheel * to ;
while ( ! TAILQ_EMPTY ( & ua - > twheel - > wheel [ w ] [ s ] ) ) {
to = TAILQ_FIRST ( & ua - > twheel - > wheel [ w ] [ s ] ) ;
TAILQ_REMOVE ( & ua - > twheel - > wheel [ w ] [ s ] , to , tqe ) ;
struct timeout_node * node = ( struct timeout_node * ) to ;
ua - > timer_free_count + + ;
memory_pool_free ( ua - > timeout_pool , node ) ;
}
}
}
timeout_heap_destroy ( ua - > timeout_heap ) ;
ua - > timeout_heap = NULL ;
struct twheel * to ;
while ( ! TAILQ_EMPTY ( & ua - > twheel - > expired ) ) {
to = TAILQ_FIRST ( & ua - > twheel - > expired ) ;
TAILQ_REMOVE ( & ua - > twheel - > expired , to , tqe ) ;
struct timeout_node * node = ( struct timeout_node * ) to ;
ua - > timer_free_count + + ;
memory_pool_free ( ua - > timeout_pool , node ) ;
}
twheels_close ( ua - > twheel ) ;
ua - > twheel = NULL ;
}
// Destroy timeout pool