Browse Source

uasync: замена timeout_heap на timing wheel (twheel)

Устранены баги потери событий:
- timeout_ms = -1 при истёкшем таймере (poll висел бесконечно)
- stale now_ms в process_timeouts (таймеры пропускались на одну итерацию)

Замена:
- timeout_heap.c/h удалены
- Добавлен twheel — tickless hierarchical timing wheel (WHEEL_BIT=6, WHEEL_NUM=4)
  на основе реализации William Ahern (MIT), адаптирован под u_malloc/u_free/DEBUG_*
- struct timeout_node: struct twheel to первым полем, убран heap_index/expiration_ms
- uasync_set_timeout: twheel_init + twheels_add(TIMEOUT_ABS)
- uasync_cancel_timeout: twheels_del + немедленный free (O(1))
- process_timeouts: twheels_update + цикл twheels_get
- get_next_timeout: twheels_timeout() — всегда актуален
- uasync_poll: исправлена логика расчёта timeout_ms
tmo
Evgeny 4 months ago
parent
commit
ce37db137b
  1. 4
      lib/Makefile.am
  2. 193
      lib/timeout_heap.c
  3. 97
      lib/timeout_heap.h
  4. 223
      lib/twheel.c
  5. 57
      lib/twheel.h
  6. 35
      lib/twheel_bitops.c
  7. 277
      lib/u_async.c
  8. 14
      lib/u_async.h
  9. 7
      src/control_server.c
  10. 1
      src/etcp.c
  11. 1
      tests/Makefile.am
  12. 8
      tests/bench_uasync_timeouts.c
  13. 1
      tests/test_u_async_comprehensive.c

4
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 \

193
lib/timeout_heap.c

@ -1,193 +0,0 @@
// timeout_heap.c
#include "timeout_heap.h"
#include "debug_config.h"
#include <stdlib.h>
#include <stdio.h> // 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;
}

97
lib/timeout_heap.h

@ -1,97 +0,0 @@
// timeout_heap.h
#ifndef TIMEOUT_HEAP_H
#define TIMEOUT_HEAP_H
#include <stdint.h> // For uint64_t
#include <stddef.h> // 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

223
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 <limits.h>
#include <stddef.h>
#include <string.h>
#include <errno.h>
#include <sys/queue.h>
#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; }

57
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 <stdbool.h>
#include <stdint.h>
#include <stdio.h>
#include <sys/queue.h>
#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

35
lib/twheel_bitops.c

@ -0,0 +1,35 @@
// twheel_bitops.c — ctz/clz for timing wheel (included by twheel.c)
#include <stdint.h>
#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 <intrin.h>
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

277
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 <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
@ -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

14
lib/u_async.h

@ -8,7 +8,6 @@
#include "platform_compat.h"
#include <stddef.h>
#include <signal.h>
#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;

7
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());

1
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;

1
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

8
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);

1
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"

Loading…
Cancel
Save