diff --git a/lib/Makefile.am b/lib/Makefile.am index eafade5d..00af0614 100644 --- a/lib/Makefile.am +++ b/lib/Makefile.am @@ -7,8 +7,8 @@ libuasync_a_SOURCES = \ ll_queue.c \ debug_config.c \ debug_config.h \ - timeout_heap.c \ - timeout_heap.h \ + twheel.c \ + twheel.h \ memory_pool.c \ memory_pool.h \ sha256.c \ diff --git a/lib/timeout_heap.c b/lib/timeout_heap.c deleted file mode 100644 index e91cf50b..00000000 --- a/lib/timeout_heap.c +++ /dev/null @@ -1,193 +0,0 @@ -// timeout_heap.c - -#include "timeout_heap.h" -#include "debug_config.h" -#include -#include // For potential error printing, optional -#include "mem.h" - -// Helper macros for 1-based indices -#define PARENT(i) ((i) / 2) -#define LEFT_CHILD(i) (2 * (i)) -#define RIGHT_CHILD(i) (2 * (i) + 1) - -TimeoutHeap *timeout_heap_create(size_t initial_capacity) { - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH1..."); - TimeoutHeap *h = u_malloc(sizeof(TimeoutHeap)); - if (!h) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH0 error..."); - return NULL; - } - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH2..."); - h->heap = u_malloc(sizeof(TimeoutEntry) * initial_capacity); - if (!h->heap) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TH1 error..."); - u_free(h); - return NULL; - } - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Creating TH3..."); - h->size = 0; - h->capacity = initial_capacity; - h->freed_count = 0; - h->user_data = NULL; - h->free_callback = NULL; - return h; -} - - -void timeout_heap_destroy(TimeoutHeap *h) { - if (!h) return; - - // Free all remaining data (deleted or not) - for (size_t i = 0; i < h->size; i++) { - if (h->free_callback) { - h->free_callback(h->user_data, h->heap[i].data); - } - } - - u_free(h->heap); - u_free(h); -} - - -static void update_index(TimeoutHeap *h, size_t idx) { - if (h->heap[idx].index_ptr) - *h->heap[idx].index_ptr = idx; -} - -static void swap_entries(TimeoutHeap *h, size_t a, size_t b) { - TimeoutEntry temp = h->heap[a]; - h->heap[a] = h->heap[b]; - h->heap[b] = temp; - update_index(h, a); - update_index(h, b); -} - -static void bubble_up(TimeoutHeap *h, size_t i) { - // i is 1-based - while (i > 1 && h->heap[PARENT(i) - 1].expiration > h->heap[i - 1].expiration) { - swap_entries(h, PARENT(i) - 1, i - 1); - i = PARENT(i); - } -} - -int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data, size_t *index_ptr) { - if (h->size == h->capacity) { - size_t new_cap = h->capacity ? h->capacity * 2 : 1; - TimeoutEntry *new_heap = u_realloc(h->heap, sizeof(TimeoutEntry) * new_cap); - if (!new_heap) return -1; // Allocation failed - h->heap = new_heap; - h->capacity = new_cap; - } - - // Insert at end (0-based) - size_t idx = h->size++; - h->heap[idx].expiration = expiration; - h->heap[idx].data = data; - h->heap[idx].index_ptr = index_ptr; - h->heap[idx].deleted = 0; - if (index_ptr) - *index_ptr = idx; - - // Bubble up (1-based) - bubble_up(h, idx + 1); - return 0; -} - -static void heapify_down(TimeoutHeap *h, size_t i) { - // i is 1-based - while (1) { - size_t smallest = i; - size_t left = LEFT_CHILD(i); - size_t right = RIGHT_CHILD(i); - - if (left <= h->size && h->heap[left - 1].expiration < h->heap[smallest - 1].expiration) { - smallest = left; - } - if (right <= h->size && h->heap[right - 1].expiration < h->heap[smallest - 1].expiration) { - smallest = right; - } - if (smallest == i) break; - - swap_entries(h, smallest - 1, i - 1); - i = smallest; - } -} - -static void remove_root(TimeoutHeap *h) { - if (h->size == 0) return; - - // Move last to root - h->heap[0] = h->heap[--h->size]; - update_index(h, 0); - - // Heapify down (1-based) - if (h->size > 0) { - heapify_down(h, 1); - } -} - -int timeout_heap_peek(TimeoutHeap *h, TimeoutEntry *out) { - if (h->size == 0) return -1; - - while (h->size > 0 && h->heap[0].deleted) { - void* data_to_free = h->heap[0].data; - remove_root(h); - if (h->free_callback) { - h->free_callback(h->user_data, data_to_free); - } - } - if (h->size == 0) return -1; - - *out = h->heap[0]; - return 0; -} - -int timeout_heap_pop(TimeoutHeap *h, TimeoutEntry *out) { - if (h->size == 0) return -1; - - while (h->size > 0 && h->heap[0].deleted) { - void* data_to_free = h->heap[0].data; - remove_root(h); - if (h->free_callback) { - h->free_callback(h->user_data, data_to_free); - } - } - if (h->size == 0) return -1; - - *out = h->heap[0]; - remove_root(h); - return 0; -} - -int timeout_heap_cancel(TimeoutHeap *h, TimeoutTime expiration, void *data) { - for (size_t i = 0; i < h->size; ++i) { - if (h->heap[i].expiration == expiration && h->heap[i].data == data) { - h->heap[i].deleted = 1; - h->heap[i].expiration = 0; - bubble_up(h, i + 1); - return 0; - } - } - return -1; // Not found -} - -int timeout_heap_cancel_at(TimeoutHeap *h, size_t index, void *data) { - if (index >= h->size || h->heap[index].data != data) return -1; - - h->heap[index].deleted = 1; - h->heap[index].expiration = 0; - bubble_up(h, index + 1); - return 0; -} - -void timeout_heap_set_free_callback(TimeoutHeap *h, void* user_data, void (*callback)(void* user_data, void* data)) { - if (!h) return; - h->user_data = user_data; - h->free_callback = callback; -} - -size_t timeout_heap_get_size(TimeoutHeap *h) { - if (!h) return 0; - return h->size; -} diff --git a/lib/timeout_heap.h b/lib/timeout_heap.h deleted file mode 100644 index 8e093f2f..00000000 --- a/lib/timeout_heap.h +++ /dev/null @@ -1,97 +0,0 @@ -// timeout_heap.h - -#ifndef TIMEOUT_HEAP_H -#define TIMEOUT_HEAP_H - -#include // For uint64_t -#include // For size_t - -typedef uint64_t TimeoutTime; // e.g., milliseconds since epoch or from now - -typedef struct { - TimeoutTime expiration; // Sort key (smaller = earlier) - void *data; // User data (e.g., callback or ID) - size_t *index_ptr; // Pointer to node's heap_index (NULL = no tracking) - int deleted; // 0 = active, 1 = deleted -} TimeoutEntry; - -typedef struct TimeoutHeap TimeoutHeap; - -struct TimeoutHeap { - TimeoutEntry *heap; // Dynamic array - size_t size; // Current number of elements - size_t capacity; // Allocated size - size_t freed_count; // Number of freed timer nodes - void* user_data; // User data for free callback - void (*free_callback)(void* user_data, void* data); // Callback to free data -}; - -/** - * Create a new timeout heap with initial capacity. - * @param initial_capacity Starting capacity (will grow as needed). - * @return Pointer to the heap, or NULL on failure. - */ -TimeoutHeap *timeout_heap_create(size_t initial_capacity); - -/** - * Destroy the timeout heap and free resources. - * @param h The heap to destroy. - */ -void timeout_heap_destroy(TimeoutHeap *h); - -/** - * Set a callback function to free data when deleted nodes are removed. - * @param h The heap. - * @param user_data User data passed to callback. - * @param callback Callback function (if NULL, data is freed with free()). - */ -void timeout_heap_set_free_callback(TimeoutHeap *h, void* user_data, void (*callback)(void* user_data, void* data)); - -/** - * Insert a new timeout into the heap. - * @param h The heap. - * @param expiration The expiration time. - * @param data User data associated with the timeout. - * @return 0 on success, -1 on allocation failure. - */ -int timeout_heap_push(TimeoutHeap *h, TimeoutTime expiration, void *data, size_t *index_ptr); - -/** - * Peek at the earliest non-deleted timeout without removing it. - * @param h The heap. - * @param out Where to store the entry. - * @return 0 on success, -1 if empty. - */ -int timeout_heap_peek(TimeoutHeap *h, TimeoutEntry *out); - -/** - * Pop the earliest non-deleted timeout from the heap. - * @param h The heap. - * @param out Where to store the entry. - * @return 0 on success, -1 if empty. - */ -int timeout_heap_pop(TimeoutHeap *h, TimeoutEntry *out); - -/** - * Cancel a timeout by matching expiration and data. - * Scans the heap linearly, so O(n) time. - * Assumes combinations are unique; cancels the first match. - * @param h The heap. - * @param expiration The expiration time to match. - * @param data The data to match. - * @return 0 if found and canceled, -1 if not found. - */ -int timeout_heap_cancel(TimeoutHeap *h, TimeoutTime expiration, void *data); - -int timeout_heap_cancel_at(TimeoutHeap *h, size_t index, void *data); - -/** - * Get the number of freed timer nodes. - * @param h The heap. - * @return Count of freed timer nodes. - */ -size_t timeout_heap_get_freed_count(TimeoutHeap *h); - -size_t timeout_heap_get_size(TimeoutHeap *h); - -#endif // TIMEOUT_HEAP_H diff --git a/lib/twheel.c b/lib/twheel.c new file mode 100644 index 00000000..ccbdbc1f --- /dev/null +++ b/lib/twheel.c @@ -0,0 +1,223 @@ +// twheel.c — Tickless hierarchical timing wheel (adapted from timeout.c) +// Original: Copyright (c) 2013-2014 William Ahern, MIT license. + +#include +#include +#include +#include +#include + +#include "twheel.h" +#include "mem.h" +#include "debug_config.h" + +#define WHEEL_BIT 6 +#define WHEEL_NUM 4 + +#define WHEEL_LEN (1U << WHEEL_BIT) +#define WHEEL_MAX (WHEEL_LEN - 1) +#define WHEEL_MASK (WHEEL_LEN - 1) + +#include "twheel_bitops.c" + +#define ctz(n) ctz64(n) +#define clz(n) clz64(n) +#define fls(n) ((int)(64 - clz64(n))) + +typedef uint64_t wheel_t; +#define WHEEL_C(n) UINT64_C(n) + +#define countof(a) (sizeof(a) / sizeof *(a)) + +#ifndef TAILQ_CONCAT +#define TAILQ_CONCAT(head1, head2, field) do { \ + if (!TAILQ_EMPTY(head2)) { \ + *(head1)->tqh_last = (head2)->tqh_first; \ + (head2)->tqh_first->field.tqe_prev = (head1)->tqh_last; \ + (head1)->tqh_last = (head2)->tqh_last; \ + TAILQ_INIT((head2)); \ + } \ +} while (0) +#endif + +#ifndef TAILQ_FOREACH_SAFE +#define TAILQ_FOREACH_SAFE(var, head, field, tvar) \ + for ((var) = TAILQ_FIRST(head); \ + (var) && ((tvar) = TAILQ_NEXT(var, field), 1); \ + (var) = (tvar)) +#endif + +#ifndef MAX +#define MAX(a,b) (((a)>(b))?(a):(b)) +#endif +#ifndef MIN +#define MIN(a,b) (((a)<(b))?(a):(b)) +#endif + +static inline wheel_t rotl(const wheel_t v, int c) { + if (!(c &= (sizeof v * CHAR_BIT - 1))) return v; + return (v << c) | (v >> (sizeof v * CHAR_BIT - c)); +} + +static inline wheel_t rotr(const wheel_t v, int c) { + if (!(c &= (sizeof v * CHAR_BIT - 1))) return v; + return (v >> c) | (v << (sizeof v * CHAR_BIT - c)); +} + +static struct twheels *twheels_init(struct twheels *T, twheel_t hz) { + for (unsigned i = 0; i < WHEEL_NUM; i++) { + for (unsigned j = 0; j < WHEEL_LEN; j++) { + TAILQ_INIT(&T->wheel[i][j]); + } + } + TAILQ_INIT(&T->expired); + for (unsigned i = 0; i < WHEEL_NUM; i++) T->pending[i] = 0; + T->curtime = 0; + T->hertz = hz ? hz : TWHEEL_mHZ; + return T; +} + +struct twheels *twheels_open(twheel_t hz) { + struct twheels *T = u_malloc(sizeof *T); + if (!T) { DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "twheels_open: malloc failed"); return NULL; } + return twheels_init(T, hz); +} + +void twheels_close(struct twheels *T) { + if (!T) return; + struct twheel_list reset; + TAILQ_INIT(&reset); + for (unsigned i = 0; i < WHEEL_NUM; i++) + for (unsigned j = 0; j < WHEEL_LEN; j++) + TAILQ_CONCAT(&reset, &T->wheel[i][j], tqe); + TAILQ_CONCAT(&reset, &T->expired, tqe); + struct twheel *to; + TAILQ_FOREACH(to, &reset, tqe) to->pending = NULL; + u_free(T); +} + +twheel_t twheels_hz(struct twheels *T) { return T->hertz; } + +void twheels_del(struct twheels *T, struct twheel *to) { + if (to->pending) { + TAILQ_REMOVE(to->pending, to, tqe); + if (to->pending != &T->expired && TAILQ_EMPTY(to->pending)) { + ptrdiff_t index = to->pending - &T->wheel[0][0]; + int w = index / WHEEL_LEN, slot = index % WHEEL_LEN; + T->pending[w] &= ~(WHEEL_C(1) << slot); + } + to->pending = NULL; + } +} + +static inline int timeout_wheel(twheel_t timeout) { + return (fls(MIN(timeout, TWHEEL_MAX)) - 1) / WHEEL_BIT; +} + +static inline int timeout_slot(int wheel, twheel_t expires) { + return WHEEL_MASK & ((expires >> (wheel * WHEEL_BIT)) - !!wheel); +} + +static void twheels_sched(struct twheels *T, struct twheel *to, twheel_t expires) { + twheels_del(T, to); + to->expires = expires; + if (expires > T->curtime) { + twheel_t rem = expires - T->curtime; + int wheel = timeout_wheel(rem); + int slot = timeout_slot(wheel, to->expires); + to->pending = &T->wheel[wheel][slot]; + TAILQ_INSERT_TAIL(to->pending, to, tqe); + T->pending[wheel] |= WHEEL_C(1) << slot; + } else { + to->pending = &T->expired; + TAILQ_INSERT_TAIL(to->pending, to, tqe); + } +} + +void twheels_add(struct twheels *T, struct twheel *to, twheel_t timeout) { + if (to->flags & TWHEEL_INT) to->interval = MAX(1, timeout); + if (to->flags & TWHEEL_ABS) twheels_sched(T, to, timeout); + else twheels_sched(T, to, T->curtime + timeout); +} + +void twheels_update(struct twheels *T, twheel_t curtime) { + twheel_t elapsed = curtime - T->curtime; + struct twheel_list todo; + TAILQ_INIT(&todo); + + for (int wheel = 0; wheel < WHEEL_NUM; wheel++) { + wheel_t pending; + if ((elapsed >> (wheel * WHEEL_BIT)) > WHEEL_MAX) { + pending = (wheel_t)~WHEEL_C(0); + } else { + wheel_t _elapsed = WHEEL_MASK & (elapsed >> (wheel * WHEEL_BIT)); + int oslot = WHEEL_MASK & (T->curtime >> (wheel * WHEEL_BIT)); + pending = rotl(((UINT64_C(1) << _elapsed) - 1), oslot); + int nslot = WHEEL_MASK & (curtime >> (wheel * WHEEL_BIT)); + pending |= rotr(rotl(((WHEEL_C(1) << _elapsed) - 1), nslot), _elapsed); + pending |= WHEEL_C(1) << nslot; + } + while (pending & T->pending[wheel]) { + int slot = ctz(pending & T->pending[wheel]); + TAILQ_CONCAT(&todo, &T->wheel[wheel][slot], tqe); + T->pending[wheel] &= ~(UINT64_C(1) << slot); + } + if (!(0x1 & pending)) break; + elapsed = MAX(elapsed, (WHEEL_LEN << (wheel * WHEEL_BIT))); + } + T->curtime = curtime; + + struct twheel *to; + while (!TAILQ_EMPTY(&todo)) { + to = TAILQ_FIRST(&todo); + TAILQ_REMOVE(&todo, to, tqe); + to->pending = NULL; + twheels_sched(T, to, to->expires); + } +} + +bool twheels_pending(struct twheels *T) { + wheel_t pending = 0; + for (int wheel = 0; wheel < WHEEL_NUM; wheel++) pending |= T->pending[wheel]; + return !!pending; +} + +bool twheels_expired(struct twheels *T) { return !TAILQ_EMPTY(&T->expired); } + +static twheel_t twheels_int(struct twheels *T) { + twheel_t timeout = ~TWHEEL_C(0), _timeout, relmask = 0; + for (int wheel = 0; wheel < WHEEL_NUM; wheel++) { + if (T->pending[wheel]) { + int slot = WHEEL_MASK & (T->curtime >> (wheel * WHEEL_BIT)); + _timeout = (ctz(rotr(T->pending[wheel], slot)) + !!wheel) << (wheel * WHEEL_BIT); + _timeout -= relmask & T->curtime; + timeout = MIN(_timeout, timeout); + } + relmask <<= WHEEL_BIT; + relmask |= WHEEL_MASK; + } + return timeout; +} + +twheel_t twheels_timeout(struct twheels *T) { + if (!TAILQ_EMPTY(&T->expired)) return 0; + return twheels_int(T); +} + +struct twheel *twheels_get(struct twheels *T) { + if (!TAILQ_EMPTY(&T->expired)) { + struct twheel *to = TAILQ_FIRST(&T->expired); + TAILQ_REMOVE(&T->expired, to, tqe); + to->pending = NULL; + return to; + } + return NULL; +} + +struct twheel *twheel_init(struct twheel *to, int flags) { + memset(to, 0, sizeof *to); + to->flags = flags; + return to; +} + +bool twheel_pending(struct twheel *to) { return to->pending != NULL; } diff --git a/lib/twheel.h b/lib/twheel.h new file mode 100644 index 00000000..33553829 --- /dev/null +++ b/lib/twheel.h @@ -0,0 +1,57 @@ +// twheel.h — Tickless hierarchical timing wheel (adapted from timeout.h) +// Original: Copyright (c) 2013-2014 William Ahern, MIT license. + +#ifndef TWHEEL_H +#define TWHEEL_H + +#include +#include +#include +#include + +#define TWHEEL_VERSION 0x20160226 + +typedef uint64_t twheel_t; +#define TWHEEL_C(n) UINT64_C(n) +#define TWHEEL_mHZ TWHEEL_C(1000) +#define TWHEEL_MAX ((TWHEEL_C(1) << 24) - 1) + +enum twheel_flags { + TWHEEL_INT = 0x01, + TWHEEL_ABS = 0x02, +}; + +#define TWHEEL_INITIALIZER(flags) { (flags) } + +struct twheel { + int flags; + twheel_t expires; + struct twheel_list *pending; + TAILQ_ENTRY(twheel) tqe; + twheel_t interval; +}; + +TAILQ_HEAD(twheel_list, twheel); + +struct twheels { + struct twheel_list wheel[4][64], expired; + uint64_t pending[4]; + twheel_t curtime; + twheel_t hertz; +}; + +struct twheel *twheel_init(struct twheel *, int flags); +bool twheel_pending(struct twheel *); + +struct twheels *twheels_open(twheel_t hz); +void twheels_close(struct twheels *); +twheel_t twheels_hz(struct twheels *); +void twheels_update(struct twheels *, twheel_t curtime); +twheel_t twheels_timeout(struct twheels *); +void twheels_add(struct twheels *, struct twheel *, twheel_t timeout); +void twheels_del(struct twheels *, struct twheel *); +struct twheel *twheels_get(struct twheels *); +bool twheels_pending(struct twheels *); +bool twheels_expired(struct twheels *); + +#endif diff --git a/lib/twheel_bitops.c b/lib/twheel_bitops.c new file mode 100644 index 00000000..58a977db --- /dev/null +++ b/lib/twheel_bitops.c @@ -0,0 +1,35 @@ +// twheel_bitops.c — ctz/clz for timing wheel (included by twheel.c) + +#include + +#if defined(__GNUC__) && !defined(TWHEEL_DISABLE_GNUC_BITOPS) +#define ctz64(n) __builtin_ctzll(n) +#define clz64(n) __builtin_clzll(n) +#if LONG_BITS == 32 +#define ctz32(n) __builtin_ctzl(n) +#define clz32(n) __builtin_clzl(n) +#else +#define ctz32(n) __builtin_ctz(n) +#define clz32(n) __builtin_clz(n) +#endif +#elif defined(_MSC_VER) && !defined(TWHEEL_DISABLE_MSVC_BITOPS) +#include +static __inline int ctz32(unsigned long val) { DWORD z=0; _BitScanForward(&z,val); return z; } +static __inline int clz32(unsigned long val) { DWORD z=0; _BitScanReverse(&z,val); return z; } +#ifdef _WIN64 +static __inline int ctz64(uint64_t val) { DWORD z=0; _BitScanForward64(&z,val); return z; } +static __inline int clz64(uint64_t val) { DWORD z=0; _BitScanReverse64(&z,val); return z; } +#else +static __inline int ctz64(uint64_t val) { uint32_t lo=(uint32_t)val,hi=(uint32_t)(val>>32); return lo?ctz32(lo):32+ctz32(hi); } +static __inline int clz64(uint64_t val) { uint32_t lo=(uint32_t)val,hi=(uint32_t)(val>>32); return hi?clz32(hi):32+clz32(lo); } +#endif +#else +#define process_(one, cz_bits, bits) if (x < ( one << (cz_bits - bits))) { rv += bits; x <<= bits; } +static inline int clz64(uint64_t x) { int rv=0; process_(UINT64_C(1),64,32); process_(UINT64_C(1),64,16); process_(UINT64_C(1),64,8); process_(UINT64_C(1),64,4); process_(UINT64_C(1),64,2); process_(UINT64_C(1),64,1); return rv; } +static inline int clz32(uint32_t x) { int rv=0; process_(UINT32_C(1),32,16); process_(UINT32_C(1),32,8); process_(UINT32_C(1),32,4); process_(UINT32_C(1),32,2); process_(UINT32_C(1),32,1); return rv; } +#undef process_ +#define process_(one, bits) if ((x & ((one << (bits))-1)) == 0) { rv += bits; x >>= bits; } +static inline int ctz64(uint64_t x) { int rv=0; process_(UINT64_C(1),32); process_(UINT64_C(1),16); process_(UINT64_C(1),8); process_(UINT64_C(1),4); process_(UINT64_C(1),2); process_(UINT64_C(1),1); return rv; } +static inline int ctz32(uint32_t x) { int rv=0; process_(UINT32_C(1),16); process_(UINT32_C(1),8); process_(UINT32_C(1),4); process_(UINT32_C(1),2); process_(UINT32_C(1),1); return rv; } +#undef process_ +#endif diff --git a/lib/u_async.c b/lib/u_async.c index 67448299..4cad19d2 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -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 #include #include @@ -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 @@ -290,16 +290,7 @@ static struct socket_node* socket_array_get_by_sock(struct socket_array* sa, soc 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 +// Simplified timeout handling without reference counting static void get_current_time(struct timeval* tv) { #ifdef _WIN32 // Для Windows используем GetTickCount или другие механизмы @@ -435,8 +426,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 +451,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 +492,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 +510,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 +540,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 +569,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; } @@ -970,65 +926,39 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events, } #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) { @@ -1330,8 +1260,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 @@ -1349,7 +1279,7 @@ struct UASYNC* uasync_create(void) { // Initialize timeout pool ua->timeout_pool = memory_pool_init(sizeof(struct timeout_node)); 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 +1294,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 +1406,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 +1469,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 diff --git a/lib/u_async.h b/lib/u_async.h index 187e12b7..a65dc712 100644 --- a/lib/u_async.h +++ b/lib/u_async.h @@ -8,7 +8,6 @@ #include "platform_compat.h" #include #include -#include "timeout_heap.h" #include "socket_compat.h" typedef void (*timeout_callback_t)(void* user_arg);// передаёт user_arg из uasync_set_timeout @@ -30,12 +29,13 @@ typedef void (*uasync_post_callback_t)(void* user_arg); #include "memory_pool.h" -struct timeout_node; // Forward declaration +struct timeout_node; +struct twheels; -struct posted_task { - uasync_post_callback_t callback; - void* arg; - struct posted_task* next; +struct posted_task { + uasync_post_callback_t callback; + void* arg; + struct posted_task* next; }; // Uasync instance structure @@ -43,7 +43,7 @@ struct UASYNC { struct memory_pool* timeout_pool; // Pool for timeout_node allocation struct timeout_node* immediate_queue_head; // FIFO queue for immediate execution struct timeout_node* immediate_queue_tail; - TimeoutHeap* timeout_heap; // Heap for timeout management + struct twheels* twheel; // Timing wheel for timeout management struct socket_array* sockets; // Array-based socket management // Debug counters for memory allocation tracking size_t timer_alloc_count; diff --git a/src/control_server.c b/src/control_server.c index d45ddea2..cf5b95ce 100644 --- a/src/control_server.c +++ b/src/control_server.c @@ -14,7 +14,6 @@ #include "pkt_normalizer.h" #include "../tools/etcpmon/etcpmon_protocol.h" #include "../lib/u_async.h" -#include "../lib/timeout_heap.h" #include "../lib/memory_pool.h" #include "../lib/debug_config.h" #include "../lib/mem.h" @@ -1028,11 +1027,7 @@ static void send_metrics(struct control_server* server, struct control_client* c } /* System resources */ - if (instance->ua && instance->ua->timeout_heap) { - rsp->etcp.active_timeouts = (uint32_t)timeout_heap_get_size(instance->ua->timeout_heap); - } else { - rsp->etcp.active_timeouts = 0; - } + rsp->etcp.active_timeouts = 0; rsp->etcp.busy_memory_blocks = (uint32_t)(u_get_allocated_count() - memory_pool_get_total_free_blocks()); diff --git a/src/etcp.c b/src/etcp.c index 3911ffa2..a7b8cb08 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -752,6 +752,7 @@ static void ack_timeout_check(struct ETCP_CONN* etcp) { // shedule timer int64_t next_timeout=timeout - elapsed; if (next_timeout<0) next_timeout=0; + if (next_timeout>timeout) next_timeout=timeout; DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] retransmission timer set for %llu units", etcp->log_name, next_timeout); etcp->retrans_timer = uasync_set_timeout(etcp->instance->ua, next_timeout+10, etcp, ack_timeout_cb, "etcp_retrans"); return; diff --git a/tests/Makefile.am b/tests/Makefile.am index 5e72c999..388f2ecf 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -42,7 +42,6 @@ check_PROGRAMS = \ test_bbr_integration \ test_intensive_memory_pool \ test_tcp_io \ - bench_timeout_heap \ bench_uasync_timeouts # Долгие тесты: запускаются только вручную, не включаются в make check diff --git a/tests/bench_uasync_timeouts.c b/tests/bench_uasync_timeouts.c index ec2b231d..78c8c325 100644 --- a/tests/bench_uasync_timeouts.c +++ b/tests/bench_uasync_timeouts.c @@ -71,7 +71,7 @@ static void benchmark_without_sockets(void) { printf("Set %d zero timeouts:\n", NUM_TIMEOUTS); printf(" Total time: %.3f ms\n", set_time); printf(" Average per timeout: %.3f us\n", (set_time * 1000.0) / NUM_TIMEOUTS); - printf(" Heap size: %zu\n\n", ua->timeout_heap->size); + printf(" Heap size: N/A (timing wheel)\n\n"); // Run mainloop 100 times clock_gettime(CLOCK_MONOTONIC, &start); @@ -84,7 +84,7 @@ static void benchmark_without_sockets(void) { printf("Run mainloop (uasync_poll) %d times:\n", NUM_TIMEOUTS); printf(" Total time: %.3f ms\n", mainloop_time); printf(" Average per iteration: %.3f us\n", (mainloop_time * 1000.0) / NUM_TIMEOUTS); - printf(" Heap size after: %zu\n\n", ua->timeout_heap->size); + printf(" Heap size after: N/A (timing wheel)\n\n"); printf("--- Results WITHOUT sockets ---\n"); printf("Set timeouts: %.3f ms (%.3f us/op)\n", set_time, (set_time * 1000.0) / NUM_TIMEOUTS); @@ -163,7 +163,7 @@ static void benchmark_with_sockets(void) { printf("Set %d zero timeouts (with %d sockets open):\n", NUM_TIMEOUTS, NUM_SOCKETS); printf(" Total time: %.3f ms\n", set_time); printf(" Average per timeout: %.3f us\n", (set_time * 1000.0) / NUM_TIMEOUTS); - printf(" Heap size: %zu\n\n", ua->timeout_heap->size); + printf(" Heap size: N/A (timing wheel)\n\n"); // Run mainloop 100 times WITH sockets clock_gettime(CLOCK_MONOTONIC, &start); @@ -176,7 +176,7 @@ static void benchmark_with_sockets(void) { printf("Run mainloop (uasync_poll) %d times (with sockets):\n", NUM_TIMEOUTS); printf(" Total time: %.3f ms\n", mainloop_time); printf(" Average per iteration: %.3f us\n", (mainloop_time * 1000.0) / NUM_TIMEOUTS); - printf(" Heap size after: %zu\n\n", ua->timeout_heap->size); + printf(" Heap size after: N/A (timing wheel)\n\n"); printf("--- Results WITH %d sockets ---\n", NUM_SOCKETS); printf("Set timeouts: %.3f ms (%.3f us/op)\n", set_time, (set_time * 1000.0) / NUM_TIMEOUTS); diff --git a/tests/test_u_async_comprehensive.c b/tests/test_u_async_comprehensive.c index b88d25b2..41c98351 100644 --- a/tests/test_u_async_comprehensive.c +++ b/tests/test_u_async_comprehensive.c @@ -10,7 +10,6 @@ #include "../lib/platform_compat.h" #include "u_async.h" -#include "timeout_heap.h" #include "debug_config.h" #include "../lib/mem.h"