Compare commits

...

10 Commits
master ... tmo

Author SHA1 Message Date
Evgeny 1f651f1df0 fix: revert etcp_link_update_inflight_lim + BBR callback in etcp_link_new 4 months ago
Evgeny cff81a70da bugfixes 4 months ago
Evgeny 4cfaaec3b0 ll_queue: inline→static inline queue_waiter_call — fix link error with -Og 4 months ago
Evgeny b5d49926a9 proxy/udp: buf 65536→1600 на стеке, UDP не нужен гигантский буфер 4 months ago
Evgeny d7c84a94e2 proxy: bounds check data_len перед memcpy в handle_data — fix buffer overflow 4 months ago
Evgeny 9af10a208f memory_pool: name field — идентификация пула в диагностике коррапшна 4 months ago
Evgeny 677a5cdcf5 log_dump: level+category, memory_pool: дамп user_data через log_dump при коррапшне 4 months ago
Evgeny 54c92f423f memory_pool: усиленная канарейка 4B + хранение alloc_loc для диагностики 4 months ago
Evgeny b37e7001b0 uasync: рефакторинг — хелперы dispatch FD/SOCK, баг process_posted_tasks, чистка 4 months ago
Evgeny ce37db137b uasync: замена timeout_heap на timing wheel (twheel) 4 months ago
  1. 4
      lib/Makefile.am
  2. 29
      lib/debug_config.c
  3. 2
      lib/debug_config.h
  4. 2
      lib/ll_queue.c
  5. 43
      lib/memory_pool.c
  6. 3
      lib/memory_pool.h
  7. 4
      lib/tcp_io.c
  8. 193
      lib/timeout_heap.c
  9. 97
      lib/timeout_heap.h
  10. 223
      lib/twheel.c
  11. 57
      lib/twheel.h
  12. 35
      lib/twheel_bitops.c
  13. 577
      lib/u_async.c
  14. 20
      lib/u_async.h
  15. 7
      src/control_server.c
  16. 2
      src/dummynet.c
  17. 5
      src/etcp.c
  18. 29
      src/etcp_connections.c
  19. 1
      src/etcp_connections.h
  20. 23
      src/lwip_tcp/lwip_tcp.c
  21. 13
      src/lwip_tcp/lwip_tcp.h
  22. 8
      src/lwip_tcp/lwip_tcp_in.c
  23. 7
      src/lwip_tcp/lwip_tcp_opts.h
  24. 4
      src/lwip_tcp/lwip_tcp_priv.h
  25. 8
      src/pkt_normalizer.c
  26. 2
      src/proxy/tcp_proxy_client.c
  27. 5
      src/proxy/tcp_proxy_server.c
  28. 2
      src/proxy/udp_proxy.c
  29. 4
      src/tun_if.c
  30. 6
      src/utun_instance.c
  31. 1
      tests/Makefile.am
  32. 6
      tests/bbr_integration/test_bbr_integration.c
  33. 8
      tests/bench_uasync_timeouts.c
  34. 6
      tests/test_etcp_congestion.c
  35. 6
      tests/test_etcp_dummynet.c
  36. 6
      tests/test_etcp_reinit_inflight.c
  37. 2
      tests/test_intensive_memory_pool.c
  38. 2
      tests/test_ll_queue.c
  39. 4
      tests/test_pkt_normalizer_standalone.c
  40. 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 \

29
lib/debug_config.c

@ -13,30 +13,29 @@
#include <stdio.h>
#include <time.h>
void log_dump(const char* prefix, const uint8_t* data, size_t len) {
void log_dump(int level, int category, const char* prefix, const uint8_t* data, size_t len) {
if (!debug_should_output(level, category)) return;
char hex_buf[513] = {'\0'};
size_t hex_len = 0;
size_t show_len = (len > 128) ? 128 : len;
for (size_t i = 0; i < show_len && hex_len < 512 - 3; i++) {
hex_len += snprintf(hex_buf + hex_len, sizeof(hex_buf) - hex_len, "%02x", data[i]);
if (i < show_len - 1 && (i + 1) % 32 == 0) { // Add space every 32 bytes
if (i < show_len - 1 && (i + 1) % 32 == 0)
hex_len += snprintf(hex_buf + hex_len, sizeof(hex_buf) - hex_len, " ");
}
}
if (len > 128) {
if (len > 128)
hex_len += snprintf(hex_buf + hex_len, sizeof(hex_buf) - hex_len, "...");
switch (level) {
case DEBUG_LEVEL_ERROR: DEBUG_ERROR(category, "%s: len=%zu hex=%s", prefix, len, hex_buf); break;
case DEBUG_LEVEL_WARN: DEBUG_WARN(category, "%s: len=%zu hex=%s", prefix, len, hex_buf); break;
case DEBUG_LEVEL_INFO: DEBUG_INFO(category, "%s: len=%zu hex=%s", prefix, len, hex_buf); break;
case DEBUG_LEVEL_DEBUG: DEBUG_DEBUG(category, "%s: len=%zu hex=%s", prefix, len, hex_buf); break;
case DEBUG_LEVEL_TRACE: DEBUG_TRACE(category, "%s: len=%zu hex=%s", prefix, len, hex_buf); break;
}
// Single-line debug output with packet info
DEBUG_INFO(DEBUG_CATEGORY_DUMP, "%s: len=%zu hex=%s", prefix, len, hex_buf);
// Additional debug info for first few bytes
// if (len >= 2) {
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "%s: first_bytes=%02x%02x last_bytes=%02x%02x",
// prefix, data[0], data[1], data[len-2], data[len-1]);
// }
}

2
lib/debug_config.h

@ -123,7 +123,7 @@ ip_str_t ip_to_str(const void *addr, int family);// big endian addr
ip_str_t sockaddr_storage_to_str(const struct sockaddr_storage *addr);// автоопределение v4 или v6
// hex dump в лог
void log_dump(const char* prefix, const uint8_t* data, size_t len);
void log_dump(int level, int category, const char* prefix, const uint8_t* data, size_t len);
/* Check if debug output should be shown for given level and category */
int debug_should_output(debug_level_t level, debug_category_t category_idx);

2
lib/ll_queue.c

@ -308,7 +308,7 @@ static void waiter_defer_cb(void* arg) {
h->defer_cb(h->defer_q, h->defer_arg);
}
inline void queue_waiter_call(struct ll_queue* q, queue_threshold_callback_fn cb, void* arg,
static inline void queue_waiter_call(struct ll_queue* q, queue_threshold_callback_fn cb, void* arg,
struct queue_waiter_handle* h) {
if (q->waiter_defer) {
h->defer_q = q; h->defer_cb = cb; h->defer_arg = arg;

43
lib/memory_pool.c

@ -5,31 +5,43 @@
#include "memory_pool.h"
#include "mem.h"
#define POOL_CANARY_OFF (pool->object_size)
#define POOL_TAG_OFF (pool->object_size + 1)
#define POOL_ALLOC_SIZE (pool->object_size + 2)
#define POOL_CANARY_VAL 0xAA
static void pool_init_tags(struct memory_pool* pool, void* obj)
// Layout: [user_data: object_size] [canary: 4 bytes] [alloc_loc: pointer] [counter: 1 byte]
#define POOL_CANARY_OFF (pool->object_size)
#define POOL_LOC_OFF (pool->object_size + 4)
#define POOL_COUNTER_OFF (pool->object_size + 4 + sizeof(const char*))
#define POOL_ALLOC_SIZE (pool->object_size + 4 + sizeof(const char*) + 1)
#define POOL_CANARY_VAL 0xDEADBEAF
static void pool_init_tags(struct memory_pool* pool, void* obj, const char* location)
{
uint8_t* canary = (uint8_t*)obj + POOL_CANARY_OFF;
uint8_t* counter = (uint8_t*)obj + POOL_TAG_OFF;
uint32_t* canary = (uint32_t*)((uint8_t*)obj + POOL_CANARY_OFF);
const char** loc = (const char**)((uint8_t*)obj + POOL_LOC_OFF);
uint8_t* counter = (uint8_t*)obj + POOL_COUNTER_OFF;
*canary = POOL_CANARY_VAL;
*loc = location;
*counter = 1;
}
static void pool_check_and_clear_tags(struct memory_pool* pool, void* obj, const char* location)
{
uint8_t* canary = (uint8_t*)obj + POOL_CANARY_OFF;
uint8_t* counter = (uint8_t*)obj + POOL_TAG_OFF;
uint32_t* canary = (uint32_t*)((uint8_t*)obj + POOL_CANARY_OFF);
const char** loc = (const char**)((uint8_t*)obj + POOL_LOC_OFF);
uint8_t* counter = (uint8_t*)obj + POOL_COUNTER_OFF;
if (*canary != POOL_CANARY_VAL) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "pool_free BUFFER OVERFLOW %p canary=0x%02x from %s! halting", obj, *canary, location);
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "pool_free BUFFER OVERFLOW pool=%p name=%s sz=%zu allocs=%zu reuse=%zu free=%d",
pool, pool->name ? pool->name : "?", pool->object_size, pool->allocations, pool->reuse_count, pool->free_count);
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, " obj=%p alloc=%s free=%s canary=0x%08x expected=0x%08x",
obj, *loc ? *loc : "(null)", location, *canary, POOL_CANARY_VAL);
if (pool->object_size)
log_dump(DEBUG_LEVEL_ERROR, DEBUG_CATEGORY_MEMORY, " user_data", obj, pool->object_size);
volatile int _halt = 1;
while (_halt) {}
}
if (*counter == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "pool_free DOUBLE FREE %p from %s! halting", obj, location);
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "pool_free DOUBLE FREE pool=%p name=%s sz=%zu allocs=%zu reuse=%zu alloc=%s free=%s",
pool, pool->name ? pool->name : "?", pool->object_size, pool->allocations, pool->reuse_count,
*loc ? *loc : "(null)", location);
volatile int _halt = 1;
while (_halt) {}
}
@ -43,7 +55,7 @@ size_t memory_pool_get_total_free_blocks(void) {
}
// Инициализировать пул памяти
struct memory_pool* memory_pool_init(size_t object_size) {
struct memory_pool* memory_pool_init(size_t object_size, const char* name) {
struct memory_pool* pool = u_calloc(1, sizeof(struct memory_pool));
if (!pool) {
return NULL;
@ -54,6 +66,7 @@ struct memory_pool* memory_pool_init(size_t object_size) {
pool->allocations = 0;
pool->reuse_count = 0;
pool->alloc_tag_counter = 1;
pool->name = name;
return pool;
}
@ -70,7 +83,7 @@ void* memory_pool_alloc_impl(struct memory_pool* pool, const char* location) {
g_total_free--;
pool->reuse_count++;
memset(obj, 0, pool->object_size);
pool_init_tags(pool, obj);
pool_init_tags(pool, obj, location);
DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "pool_alloc reused: %s, remaining=%zu", location, pool->free_count);
return obj;
}
@ -78,7 +91,7 @@ void* memory_pool_alloc_impl(struct memory_pool* pool, const char* location) {
pool->allocations++;
DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "pool_alloc: %s, total_allocs=%zu", location, pool->allocations);
void* obj = u_calloc_impl(1, POOL_ALLOC_SIZE, location);
if (obj) pool_init_tags(pool, obj);
if (obj) pool_init_tags(pool, obj, location);
return obj;
}

3
lib/memory_pool.h

@ -18,11 +18,12 @@ struct memory_pool {
size_t allocations; // Общее количество аллокаций (включая новые malloc)
size_t reuse_count; // Количество повторных использований из пула
uint8_t alloc_tag_counter; // монотонный счётчик для детекции double-free (1 байт)
const char* name; // Имя пула для диагностики (напр. "data_pool", "pkt_pool")
};
// сам пул:
// если используем с элементами ll_queue то не забываем что к object_size надо прибавить sizeof(struct ll_entry)
struct memory_pool* memory_pool_init(size_t object_size);
struct memory_pool* memory_pool_init(size_t object_size, const char* name);
void memory_pool_destroy(struct memory_pool* pool);
// элементы пула:

4
lib/tcp_io.c

@ -46,10 +46,10 @@ struct tcp_conn* tcp_conn_create(
tc->on_error = on_error;
tc->arg = arg;
tc->entry_pool = memory_pool_init(sizeof(struct ll_entry));
tc->entry_pool = memory_pool_init(sizeof(struct ll_entry), "entry_pool");
{
size_t ds = entry_data_size > write_chunk_size ? entry_data_size : write_chunk_size;
tc->data_pool = memory_pool_init(ds);
tc->data_pool = memory_pool_init(ds, "data_pool");
}
if (!tc->entry_pool || !tc->data_pool) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_conn_create: memory_pool_init failed");

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

577
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
@ -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_TUN, "POSTed task get");
DEBUG_DEBUG(DEBUG_CATEGORY_UASYNC, "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_read, 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

20
lib/u_async.h

@ -6,9 +6,9 @@
#define UASYNC_H
#include "platform_compat.h"
#include "memory_pool.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
@ -26,16 +26,15 @@ typedef int err_t;
#define SOCKET_NODE_TYPE_FD 0 // Regular file descriptor (pipe, file)
#define SOCKET_NODE_TYPE_SOCK 1 // Socket (socket_t)
typedef void (*uasync_post_callback_t)(void* user_arg);
#include "memory_pool.h"
typedef void (*uasync_post_callback_t)(void* user_arg);
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 +42,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;
@ -74,7 +73,6 @@ struct UASYNC {
// Type definitions
typedef struct UASYNC uasync_t;
typedef struct UASYNC UASYNC_t;
// Instance API - основной API для работы с uasync
struct UASYNC* uasync_create(void);

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

2
src/dummynet.c

@ -355,7 +355,7 @@ struct dummynet* dummynet_create(struct UASYNC* ua, const char* bind_ip, uint16_
dn->listen_port = listen_port;
/* Создаём пул для пакетов */
dn->pkt_pool = memory_pool_init(sizeof(struct ll_entry) + sizeof(struct dummynet_pkt));
dn->pkt_pool = memory_pool_init(sizeof(struct ll_entry) + sizeof(struct dummynet_pkt), "pkt_pool");
if (!dn->pkt_pool) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to create packet pool");
u_free(dn);

5
src/etcp.c

@ -205,8 +205,8 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n
etcp->input_wait_ack = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "input_wait_ack"); // Hash for wait_ack
etcp->recv_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "recv_q"); // Hash for recv_q
etcp->ack_q = queue_new(instance->ua, INFLIGHT_INITIAL_HASH_SIZE, 0, 4, "ack_q");
etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET));
etcp->io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT));
etcp->inflight_pool = memory_pool_init(sizeof(struct INFLIGHT_PACKET), "inflight_pool");
etcp->io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT), "io_pool");
etcp->optimal_inflight=100000;
etcp->initialized=0;
etcp->links_up=0;
@ -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;

29
src/etcp_connections.c

@ -746,8 +746,6 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) {
}
static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd);
struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!remote_addr) return NULL;
@ -833,9 +831,12 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
while (l && l->next) l=l->next;
if (l) l->next = link; else etcp->links = link;
etcp_link_update_inflight_lim(link, link->mtu * 4);
link->bbr->on_cwnd_update = bbr_cwnd_updated;
link->bbr->cwnd_update_ctx = link;
link->inflight_lim_bytes = link->mtu * 4; // BBR init_cwnd (~4 packets)
// пересчитать connection-level optimal_inflight
{ uint32_t sum = 0;
for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes;
etcp->optimal_inflight = sum; }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "NEW link initialized on etcp=[%s] link=%p socket=%s id=%d is_server=%d mtu=%d", etcp->log_name, link, conn->name, link->local_link_id, link->is_server, link->mtu);
@ -847,18 +848,6 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
return link;
}
static void bbr_cwnd_updated(void* ctx, uint32_t new_cwnd) {
etcp_link_update_inflight_lim((struct ETCP_LINK*)ctx, new_cwnd);
}
void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) {
link->inflight_lim_bytes = new_lim;
struct ETCP_CONN* etcp = link->etcp;
uint32_t sum = 0;
for (struct ETCP_LINK* tl = etcp->links; tl; tl = tl->next) sum += tl->inflight_lim_bytes;
etcp->optimal_inflight = sum;
}
void etcp_link_close(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!link) return;
@ -940,7 +929,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
dgram->timestamp=get_current_timestamp();
dgram->link->total_encrypted += dgram->data_len;
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump("Before encryption", dgram->data, dgram->data_len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "Before encryption", dgram->data, dgram->data_len);
sc_encrypt(sc, (uint8_t*)&dgram->timestamp, 3 + len, enc_buf, &enc_buf_len);
if (enc_buf_len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "eencryption failed for node %016llx", (unsigned long long)dgram->link->etcp->instance->node_id);
@ -950,7 +939,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
dgram->link->send_errors++; errcode=3; goto es_err; }
memcpy(enc_buf+enc_buf_len, dgram->data+len, dgram->noencrypt_len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump("Encrypted", enc_buf, enc_buf_len + dgram->noencrypt_len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "Encrypted", enc_buf, enc_buf_len + dgram->noencrypt_len);
struct sockaddr_storage* addr=&dgram->link->remote_addr;
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6);
@ -1180,7 +1169,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
}
// DUMP: Show received packet content
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump("RECV in:", data, recv_len); // link unknown at this point
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "RECV in:", data, recv_len);
struct ETCP_DGRAM* pkt = memory_pool_alloc(e_sock->instance->pkt_pool);
if (!pkt) return;

1
src/etcp_connections.h

@ -262,7 +262,6 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn);
// connection functions
// создает новый канал связи для etcp подключения (ETCP_CONN)
struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server);
void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim);
void etcp_link_close(struct ETCP_LINK* link);
//int etcp_input_cbk(struct packet_buffer* pkt, struct ETCP_SOCKET* conn);// получает расшифрованный пакет
int etcp_encrypt_send(struct ETCP_DGRAM* dgram);// зашифровывает и отправляет пакет

23
src/lwip_tcp/lwip_tcp.c

@ -95,15 +95,15 @@ struct lwip_tcp_ctx *lwip_tcp_init(struct UASYNC *ua, tcp_output_fn output, void
ctx->ua = ua;
ctx->output = output;
ctx->output_arg = output_arg;
ctx->pcb_pool = memory_pool_init(sizeof(struct tcp_pcb));
ctx->pcb_listen_pool = memory_pool_init(sizeof(struct tcp_pcb_listen));
ctx->seg_pool = memory_pool_init(sizeof(struct tcp_seg));
ctx->pcb_pool = memory_pool_init(sizeof(struct tcp_pcb), "pcb_pool");
ctx->pcb_listen_pool = memory_pool_init(sizeof(struct tcp_pcb_listen), "pcb_listen_pool");
ctx->seg_pool = memory_pool_init(sizeof(struct tcp_seg), "seg_pool");
if (!ctx->pcb_pool || !ctx->pcb_listen_pool || !ctx->seg_pool) {
lwip_tcp_destroy(ctx);
return NULL;
}
ctx->iss_seed = (uint16_t)(get_time_tb() & 0xFFFF);
ctx->tmr_interval_ms = TCP_TMR_INTERVAL;
ctx->tmr_interval_ms = TCP_TMR_INTERVAL / 4;
ctx->rto_min_ms = 3000;
ctx->rto_max_ms = 0;
ctx->timer = uasync_set_timeout(ua, ctx->tmr_interval_ms * 10, ctx, tcp_tmr_cb, "lwip_tcp_tmr");
@ -749,7 +749,6 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx)
tcp_pcb_purge(pcb);
if (prev != NULL) {
prev->next = pcb->next;
prev->next_owner = PCB_NEXT_SLOWTMR;
} else {
ctx->active_pcbs = pcb->next;
}
@ -808,18 +807,16 @@ void tcp_slowtmr(struct lwip_tcp_ctx *ctx)
tcp_pcb_purge(pcb);
if (prev != NULL) {
prev->next = pcb->next;
prev->next_owner = PCB_NEXT_SLOWTMR;
} else {
ctx->tw_pcbs = pcb->next;
}
pcb2 = pcb;
pcb = pcb->next;
int self_loop = (pcb == pcb2);
uint8_t sl_owner = pcb2->next_owner;
tcp_free(pcb2);
if (self_loop) {
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_slowtmr: freed pcb=%p state=TIME_WAIT port=%u next_owner=%d, clearing tw_pcbs",
(void*)pcb2, pcb2->local_port, sl_owner);
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_slowtmr: freed pcb=%p state=TIME_WAIT port=%u, clearing tw_pcbs",
(void*)pcb2, pcb2->local_port);
ctx->tw_pcbs = NULL;
break;
}
@ -1050,8 +1047,8 @@ static void tcp_kill_timewait(struct lwip_tcp_ctx *ctx)
}
if (inactive != NULL) {
if (inactive->next == inactive) {
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_kill_timewait: aborting pcb=%p next_owner=%d, clearing tw_pcbs",
(void*)inactive, inactive->next_owner);
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in tcp_kill_timewait: aborting pcb=%p, clearing tw_pcbs",
(void*)inactive);
ctx->tw_pcbs = NULL;
tcp_free(inactive);
} else {
@ -1115,8 +1112,8 @@ struct tcp_pcb *tcp_alloc(struct lwip_tcp_ctx *ctx, uint8_t prio)
pcb->rcv_wnd = pcb->rcv_ann_wnd = TCPWND16(TCP_WND);
pcb->ttl = 64;
pcb->mss = INITIAL_MSS;
pcb->rto = (int16_t)(ctx->rto_min_ms / TCP_SLOW_INTERVAL);
pcb->sv = (int16_t)(ctx->rto_min_ms / TCP_SLOW_INTERVAL);
pcb->rto = (int16_t)(TCP_RTO_MIN_MS / TCP_SLOW_INTERVAL);
pcb->sv = (int16_t)(TCP_RTO_MIN_MS / TCP_SLOW_INTERVAL);
pcb->rtime = -1;
pcb->cwnd = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "tcp_alloc: cwnd=%u snd_buf=%u mss=%u snd_wnd=%u", pcb->cwnd, pcb->snd_buf, pcb->mss, pcb->snd_wnd);

13
src/lwip_tcp/lwip_tcp.h

@ -67,15 +67,6 @@ enum tcp_err_enum {
LERR_CLSD = -15,
LERR_ARG = -16
};
// кто последним записал pcb->next (диагностика зацикливания tw_pcbs)
enum pcb_next_owner {
PCB_NEXT_NONE = 0,
PCB_NEXT_REG = 1,
PCB_NEXT_RMV = 2,
PCB_NEXT_SLOWTMR = 3,
PCB_NEXT_INPUT = 4,
};
typedef int err_t;
// Forward declaration for callbacks
@ -175,9 +166,6 @@ struct tcp_pcb {
tcp_connected_fn connected;
tcp_poll_fn poll;
tcp_err_fn errf;
uint8_t next_owner; // who last wrote pcb->next (enum pcb_next_owner)
uint32_t keep_idle;
uint8_t persist_cnt;
uint8_t persist_backoff;
@ -193,7 +181,6 @@ struct tcp_pcb_listen {
uint8_t prio;
uint16_t local_port;
uint32_t local_ip;
uint8_t next_owner;
tcp_accept_fn accept;
};

8
src/lwip_tcp/lwip_tcp_in.c

@ -198,9 +198,7 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p,
pcb->local_ip == dst_ip) {
if (prev != NULL) {
prev->next = pcb->next;
prev->next_owner = PCB_NEXT_INPUT;
pcb->next = ctx->active_pcbs;
pcb->next_owner = PCB_NEXT_INPUT;
ctx->active_pcbs = pcb;
}
break;
@ -225,8 +223,8 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p,
pcb->remote_ip == src_ip &&
pcb->local_ip == dst_ip) {
if (pcb->next == pcb) {
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in lwip_tcp_input: pcb=%p sport=%u dport=%u next_owner=%d, aborting",
(void*)pcb, sport, dport, pcb->next_owner);
DEBUG_ERROR(DEBUG_CATEGORY_ALL, "TW_PCBS SELF-LOOP in lwip_tcp_input: pcb=%p sport=%u dport=%u, aborting",
(void*)pcb, sport, dport);
tcp_abort(pcb);
} else {
tcp_timewait_input(pcb);
@ -249,9 +247,7 @@ void lwip_tcp_input(struct lwip_tcp_ctx *ctx, struct pbuf *p,
if (lpcb != NULL) {
if (prev != NULL) {
((struct tcp_pcb_listen *)prev)->next = lpcb->next;
((struct tcp_pcb_listen *)prev)->next_owner = PCB_NEXT_INPUT;
lpcb->next = ctx->listen_pcbs;
lpcb->next_owner = PCB_NEXT_INPUT;
ctx->listen_pcbs = (struct tcp_pcb *)lpcb;
}
tcp_listen_input(lpcb);

7
src/lwip_tcp/lwip_tcp_opts.h

@ -14,9 +14,10 @@
#define TCP_FAST_INTERVAL TCP_TMR_INTERVAL
#define TCP_SLOW_INTERVAL (2 * TCP_TMR_INTERVAL)
#define TCP_FIN_WAIT_TIMEOUT 6000 // ms
#define TCP_SYN_RCVD_TIMEOUT 6000 // ms
#define TCP_MSL 25000 // ms (2*MSL = 50s)
#define TCP_RTO_MIN_MS 3000 // ms, initial RTO
#define TCP_FIN_WAIT_TIMEOUT 20000 // ms
#define TCP_SYN_RCVD_TIMEOUT 20000 // ms
#define TCP_MSL 60000 // ms (2*MSL = 120s)
#define TCP_OOSEQ_TIMEOUT 6 // x RTO
#define TCP_KEEPIDLE_DEFAULT 7200000 // ms (unused, no keepalive)

4
src/lwip_tcp/lwip_tcp_priv.h

@ -236,15 +236,13 @@ void lwip_tcp_stats_clear(struct lwip_tcp_ctx *ctx);
// PCB list management
#define TCP_REG(pcbs, npcb) do { \
(npcb)->next_owner = PCB_NEXT_REG; \
(npcb)->next = *(pcbs); *(pcbs) = (npcb); \
} while(0)
#define TCP_RMV(pcbs, npcb) do { \
if(*(pcbs) == (npcb)) { *(pcbs) = (*pcbs)->next; } \
else { struct tcp_pcb *_tmp; for(_tmp = *(pcbs); _tmp != NULL; _tmp = _tmp->next) { if(_tmp->next == (npcb)) { _tmp->next_owner = PCB_NEXT_SLOWTMR; _tmp->next = (npcb)->next; break; } } } \
else { struct tcp_pcb *_tmp; for(_tmp = *(pcbs); _tmp != NULL; _tmp = _tmp->next) { if(_tmp->next == (npcb)) { _tmp->next = (npcb)->next; break; } } } \
(npcb)->next = NULL; \
(npcb)->next_owner = PCB_NEXT_RMV; \
} while(0)
#define TCP_REG_ACTIVE(ctx, npcb) TCP_REG(&(ctx)->active_pcbs, npcb)

8
src/pkt_normalizer.c

@ -241,7 +241,7 @@ static void pn_send_to_etcp(struct PKTNORM* pn) {
frag->ll.memlen = pn->etcp->instance->data_pool->object_size;
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn->etcp: size=%d memlen=%d frag_size=%d", frag->ll.len, frag->ll.memlen, pn->frag_size);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->ETCP", pn->data, frag->ll.len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "NORM->ETCP", pn->data, frag->ll.len);
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
@ -286,7 +286,7 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
return;
}
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("->NORM", in_dgram->dgram, in_dgram->len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "->NORM", in_dgram->dgram, in_dgram->len);
pn_buf_renew(pn);
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: new pkt hdrpos=%d",pn->data_ptr);
@ -350,7 +350,7 @@ static void pn_unpacker_cb(struct ll_queue* q, void* arg) {
uint8_t* payload = frag->ll.dgram;
uint16_t len = frag->ll.len;
uint16_t ptr = 0;
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("ETCP->NORM", payload, len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "ETCP->NORM", payload, len);
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacking fragment len=%d", len);
@ -395,7 +395,7 @@ static void pn_unpacker_cb(struct ll_queue* q, void* arg) {
if (pn->recvpart->len == pn->recvpart->memlen) {
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacked dgram (size=%d)", pn->recvpart->len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->", pn->recvpart->dgram, pn->recvpart->len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "NORM->", pn->recvpart->dgram, pn->recvpart->len);
queue_data_put(pn->output, pn->recvpart);
pn->out_total_pkts++;
pn->out_total_bytes += pn->recvpart->len;

2
src/proxy/tcp_proxy_client.c

@ -564,7 +564,7 @@ struct tcp_proxy_client* tcp_proxy_client_create(struct UTUN_INSTANCE* inst, str
p->via_node_id = via_node_id;
p->mappings = mappings; p->mapping_count = mapping_count;
p->entry_pool = memory_pool_init(sizeof(struct ll_entry));
p->entry_pool = memory_pool_init(sizeof(struct ll_entry), "entry_pool");
if (!p->entry_pool) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_client_create: memory_pool_init failed"); u_free(p); return NULL; }
p->tun = tun_init_nat(ua, tun_name, tun_ip, mtu > 0 ? mtu : 1500, test_mode);

5
src/proxy/tcp_proxy_server.c

@ -325,6 +325,11 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c
}
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
if (data_len > 0) {
if (data_len > rc->tc->data_pool->object_size) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: data_len=%zu > pool_sz=%zu, dropping sid=%08x",
data_len, rc->tc->data_pool->object_size, stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return -1;
}
struct ll_entry* e = queue_entry_new_from_pool(rc->tc->entry_pool);
uint8_t* buf = memory_pool_alloc(rc->tc->data_pool);
if (e && buf) {

2
src/proxy/udp_proxy.c

@ -40,7 +40,7 @@ static struct udp_flow* flow_find(struct udp_flow* head, uint64_t client_node_id
static void flow_read_cb(socket_t sock, void* arg) {
(void)sock; struct udp_flow* f = (struct udp_flow*)arg;
if (!f || f->sock == SOCKET_INVALID || !g_udp_ctx) return;
uint8_t buf[65536];
uint8_t buf[1600];
struct sockaddr_in from; socklen_t flen = sizeof(from);
ssize_t n = recvfrom(f->sock, buf, sizeof(buf), 0, (struct sockaddr*)&from, &flen);
if (n <= 0) return;

4
src/tun_if.c

@ -161,7 +161,7 @@ struct tun_if* tun_init(struct UASYNC* ua, struct utun_config* config)
tun->test_fd = fds[1];
}
tun->pool = memory_pool_init(sizeof(struct ll_entry));
tun->pool = memory_pool_init(sizeof(struct ll_entry), "tun_pool");
if (!tun->pool) goto fail;
tun->output_queue = queue_new(ua, 0, 0, 0, "TUN output");
@ -253,7 +253,7 @@ struct tun_if* tun_init_nat(struct UASYNC* ua, const char* ifname, const char* i
tun->test_fd = fds[1];
}
tun->pool = memory_pool_init(sizeof(struct ll_entry));
tun->pool = memory_pool_init(sizeof(struct ll_entry), "tun_pool");
if (!tun->pool) goto fail2;
tun->output_queue = queue_new(ua, 0, 0, 0, "NAT TUN output");

6
src/utun_instance.c

@ -83,9 +83,9 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
}
// Create memory pools
instance->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET));
instance->data_pool = memory_pool_init(PACKET_DATA_SIZE);
instance->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
instance->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
instance->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
instance->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
// Create routing module
if (routing_create(instance) != 0) {

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

6
tests/bbr_integration/test_bbr_integration.c

@ -111,9 +111,9 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
inst->ua = u;
inst->node_id = node_id;
if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; }
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET));
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE);
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; }
struct utun_config* cfg = u_calloc(1, sizeof(*cfg));
if (!cfg) { u_free(inst); return NULL; }

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

6
tests/test_etcp_congestion.c

@ -102,9 +102,9 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
inst->ua = u;
inst->node_id = node_id;
if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; }
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET));
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE);
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; }
struct utun_config* cfg = u_calloc(1, sizeof(*cfg));
if (!cfg) { u_free(inst); return NULL; }

6
tests/test_etcp_dummynet.c

@ -111,9 +111,9 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
return NULL;
}
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET));
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE);
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) {
u_free(inst);

6
tests/test_etcp_reinit_inflight.c

@ -53,9 +53,9 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
inst->ua = u;
inst->node_id = node_id;
if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; }
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET));
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE);
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; }
struct utun_config* cfg = u_calloc(1, sizeof(*cfg));
if (!cfg) { u_free(inst); return NULL; }

2
tests/test_intensive_memory_pool.c

@ -65,7 +65,7 @@ static double test_with_pools(int iterations) {
struct ll_queue* queue = queue_new(ua, 0, 0, 0,"q2"); // С пулами
// Создать пул памяти для данных
struct memory_pool* pool = memory_pool_init(sizeof(struct ll_entry) + 64);
struct memory_pool* pool = memory_pool_init(sizeof(struct ll_entry) + 64, "test_pool");
if (!pool) {
queue_free(queue);
uasync_destroy(ua, 0);

2
tests/test_ll_queue.c

@ -354,7 +354,7 @@ static void test_limits_hash(void) {
static void test_pool(void) {
TEST("memory_pool integration + reuse");
struct UASYNC *ua = uasync_create();
struct memory_pool *pool = memory_pool_init(sizeof(test_data_t));
struct memory_pool *pool = memory_pool_init(sizeof(test_data_t), "test_data_pool");
struct ll_queue *q = queue_new(ua, 0, 0, 0, "q8");
size_t alloc1 = 0, reuse1 = 0;

4
tests/test_pkt_normalizer_standalone.c

@ -185,7 +185,7 @@ static int init_mock_etcp(void) {
mock_etcp.instance = (struct UTUN_INSTANCE*)&mock_instance;
// Create io_pool for ETCP_FRAGMENT allocation
mock_etcp.io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT));
mock_etcp.io_pool = memory_pool_init(sizeof(struct ETCP_FRAGMENT), "io_pool");
if (!mock_etcp.io_pool) {
printf("Failed to create io_pool\n");
return -1;
@ -271,7 +271,7 @@ int main() {
}
// Initialize memory pool
mock_instance.data_pool = memory_pool_init(MTU_SIZE);
mock_instance.data_pool = memory_pool_init(MTU_SIZE, "data_pool");
if (!mock_instance.data_pool) {
printf("Failed to create memory pool\n");

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