From d0af315ff629fb7524253ce7973985c1fae45cc6 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Tue, 12 May 2026 13:42:16 +0300 Subject: [PATCH] fix uasync set_timeout(0) ordering, dummynet shaper initial delay, add tests - uasync_set_timeout(0): use FIFO immediate_queue instead of timeout_heap to guarantee execution AFTER all currently-expired heap timers - uasync_cancel_timeout: search immediate_queue for timeout=0 timers - dummynet: start initial shaper with bandwidth-based delay instead of 0 (prevents queue from never filling up when packets arrive sequentially) - test_dummynet: rename get_time_us -> now_us (conflict with u_async.h) - test_u_async_comprehensive: fix immediate_timeouts, add timeout_ordering test --- lib/u_async.c | 50 +++++++++++++++++++--------- src/dummynet.c | 11 +++++-- tests/test_dummynet.c | 14 ++++---- tests/test_u_async_comprehensive.c | 52 ++++++++++++++++++++++++++---- 4 files changed, 96 insertions(+), 31 deletions(-) diff --git a/lib/u_async.c b/lib/u_async.c index 7af7f93c..b784b9a1 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -529,23 +529,32 @@ void* uasync_set_timeout(struct UASYNC* ua, int timeout_tb, void* arg, timeout_c } node->arg = arg; node->callback = callback; - node->ua = ua; - - // Calculate expiration time in milliseconds - struct timeval now; - get_current_time(&now); - timeval_add_tb(&now, timeout_tb); - node->expiration_ms = timeval_to_ms(&now); - + node->ua = ua; + + if (timeout_tb == 0) { + node->expiration_ms = 0; + node->next = NULL; + if (ua->immediate_queue_tail) ua->immediate_queue_tail->next = node; + else ua->immediate_queue_head = node; + ua->immediate_queue_tail = node; + return node; + } + + // Calculate expiration time in milliseconds + struct timeval now; + get_current_time(&now); + timeval_add_tb(&now, timeout_tb); + node->expiration_ms = timeval_to_ms(&now); + // Add to heap if (timeout_heap_push(ua->timeout_heap, node->expiration_ms, node) != 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; + } + + return node; } // Immediate execution in next mainloop (FIFO order) @@ -611,10 +620,21 @@ err_t uasync_cancel_timeout(struct UASYNC* ua, void* t_id) { return ERR_OK; } - 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 not found in heap, it may have already expired or been invalid + 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); + + // Not in heap — try immediate_queue (for timeout=0 timers placed via FIFO) + { + struct timeout_node* cur = ua->immediate_queue_head; + while (cur) { + if (cur == node) { + cur->callback = NULL; + return ERR_OK; + } + cur = cur->next; + } + } + return ERR_FAIL; } diff --git a/src/dummynet.c b/src/dummynet.c index 0e8d7870..b6c69317 100644 --- a/src/dummynet.c +++ b/src/dummynet.c @@ -131,7 +131,6 @@ static void dummynet_process_queue(struct dummynet* dn, int dir_idx) { } struct dummynet_pkt* pkt = (struct dummynet_pkt*)entry->data; - ssize_t n = socket_sendto(dn->sock, pkt->data, pkt->len, (struct sockaddr*)&dir->dest_addr, dir->dest_len); if (n < 0) { @@ -327,7 +326,15 @@ static void dummynet_delay_callback(void* user_arg) { if (sarg) { sarg->dn = dn; sarg->dir_idx = dir_idx; - dir->shaper_timer = uasync_set_timeout(dn->ua, 0, sarg, dummynet_shaper_callback, "dummynet_shaper0"); + int shaper_delay_tb = 0; + if (dir->bandwidth_kbps > 0) { + uint64_t delay_calc = (uint64_t)pkt->len * 80ULL / (uint64_t)dir->bandwidth_kbps; + shaper_delay_tb = (int)delay_calc; + if (shaper_delay_tb == 0) shaper_delay_tb = 1; + } + dir->shaper_timer = uasync_set_timeout(dn->ua, shaper_delay_tb, sarg, dummynet_shaper_callback, "dummynet_shaper0"); + if (dir->shaper_timer) dir->shaper_timer_set++; + else { u_free(sarg); dir->shaper_pending = 0; } } else { dir->shaper_pending = 0; } diff --git a/tests/test_dummynet.c b/tests/test_dummynet.c index 4b7b1713..c4907ff3 100644 --- a/tests/test_dummynet.c +++ b/tests/test_dummynet.c @@ -60,7 +60,7 @@ struct test_state { int passed; }; -static inline uint64_t get_time_us(void) { +static inline uint64_t now_us(void) { struct timeval tv; gettimeofday(&tv, NULL); return (uint64_t)tv.tv_sec * 1000000ULL + tv.tv_usec; @@ -75,7 +75,7 @@ static void send_one_packet(struct test_state* st, uint32_t seq) { uint16_t data_len = MIN_PKT_SIZE + (rand() % (max_data_len - MIN_PKT_SIZE + 1)); pkt->seq_num = seq; - pkt->send_time_us = get_time_us(); + pkt->send_time_us = now_us(); pkt->data_len = data_len; /* Заполняем данные */ @@ -95,7 +95,7 @@ static void send_one_packet(struct test_state* st, uint32_t seq) { /* B -> dummynet (backward) */ pkt->seq_num = seq + 10000; /* Отличаем пакеты от B */ - pkt->send_time_us = get_time_us(); + pkt->send_time_us = now_us(); socket_sendto(st->sock_b, buf, sizeof(*pkt) + data_len, (struct sockaddr*)&dest, sizeof(dest)); } @@ -136,7 +136,7 @@ static void recv_callback_a(socket_t sock, void* arg) { if (n < (ssize_t)sizeof(struct test_pkt)) return; struct test_pkt* pkt = (struct test_pkt*)buf; - uint64_t now = get_time_us(); + uint64_t now = now_us(); uint64_t latency = now - pkt->send_time_us; st->recv_a++; @@ -155,7 +155,7 @@ static void recv_callback_b(socket_t sock, void* arg) { if (n < (ssize_t)sizeof(struct test_pkt)) return; struct test_pkt* pkt = (struct test_pkt*)buf; - uint64_t now = get_time_us(); + uint64_t now = now_us(); uint64_t latency = now - pkt->send_time_us; st->recv_b++; @@ -262,7 +262,7 @@ static int run_scenario(const char* name, st.send_timer = uasync_set_timeout(st.ua, 1, &st, send_timer_callback, "test_dummynet"); /* Умное ожидание завершения: ранний выход при готовности + watchdog */ - uint64_t start_time = get_time_us(); + uint64_t start_time = now_us(); uint64_t watchdog_start = 0; /* Стартует после отправки всех пакетов */ int watchdog_active = 0; const uint64_t WATCHDOG_US = 3000000ULL; /* 3 секунды максимум после отправки */ @@ -273,7 +273,7 @@ static int run_scenario(const char* name, while (1) { uasync_poll(st.ua, 1); /* 0.1ms timeout - чаще для обработки таймеров */ - uint64_t now = get_time_us(); + uint64_t now = now_us(); /* Активируем watchdog когда все пакеты отправлены */ if (!watchdog_active && st.next_seq >= PKT_COUNT) { diff --git a/tests/test_u_async_comprehensive.c b/tests/test_u_async_comprehensive.c index 92a8d0ce..9729bd39 100644 --- a/tests/test_u_async_comprehensive.c +++ b/tests/test_u_async_comprehensive.c @@ -242,14 +242,9 @@ static void test_immediate_timeouts(void) { ASSERT_NOT_NULL(timer, "Failed to set immediate timer"); } - /* Immediate timeouts should fire during next poll */ + /* Immediate (timeout=0) callbacks use FIFO queue — all fire in first poll */ ASSERT_EQ(ctx.callback_count, 0, "Callbacks fired too early"); - - /* Note: Library processes only one timeout per poll, so we need multiple polls */ - for (int i = 0; i < 5; i++) { - uasync_poll(ua, 1); /* Minimal poll */ - } - + uasync_poll(ua, 1); ASSERT_EQ(ctx.callback_count, ctx.expected_count, "Immediate timeouts didn't fire correctly"); /* Check that exactly 5 new immediate timeouts were recorded */ @@ -260,6 +255,48 @@ static void test_immediate_timeouts(void) { TEST_PASS(); } +/* Test 3b: Timeout ordering — timeout=0 set during callback fires AFTER all expired heap timers */ +struct ordering_ctx { + int seq[8]; + int idx; + uasync_t* ua; +}; + +static void order_immediate_cb(void* arg); + +static void order_delayed_cb(void* arg) { + struct ordering_ctx* ctx = (struct ordering_ctx*)arg; + ctx->seq[ctx->idx++] = 1; + if (ctx->idx == 1) uasync_set_timeout(ctx->ua, 0, ctx, order_immediate_cb, "test_order_imm"); +} + +static void order_immediate_cb(void* arg) { + struct ordering_ctx* ctx = (struct ordering_ctx*)arg; + ctx->seq[ctx->idx++] = 2; +} + +static void test_timeout_ordering(void) { + TEST_START("Timeout ordering: immediate after delayed"); + + uasync_t* ua = uasync_create(); + ASSERT_NOT_NULL(ua, "Failed to create uasync instance"); + + struct ordering_ctx ctx = {{0}}; + ctx.ua = ua; + + for (int i = 0; i < 5; i++) + ASSERT_NOT_NULL(uasync_set_timeout(ua, 10, &ctx, order_delayed_cb, "test_order_del"), "Failed to set timer"); + + for (int i = 0; i < 20 && ctx.idx < 6; i++) uasync_poll(ua, 100); /* 10ms — enough for 1ms delay timers */ + + ASSERT_EQ(ctx.idx, 6, "Not all callbacks fired"); + for (int i = 0; i < 5; i++) ASSERT_EQ(ctx.seq[i], 1, "Delayed callback out of order"); + ASSERT_EQ(ctx.seq[5], 2, "Immediate callback fired before all delayed"); + + uasync_destroy(ua, 0); + TEST_PASS(); +} + /* Test 4: Memory leak detection */ static void test_memory_leak_detection(void) { TEST_START("Memory leak detection"); @@ -502,6 +539,7 @@ int main(void) { test_basic_timers(); test_timer_cancellation_races(); test_immediate_timeouts(); + test_timeout_ordering(); test_memory_leak_detection(); test_socket_management(); test_error_handling();