Browse Source

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
congestion
Evgeny 5 months ago
parent
commit
d0af315ff6
  1. 50
      lib/u_async.c
  2. 11
      src/dummynet.c
  3. 14
      tests/test_dummynet.c
  4. 52
      tests/test_u_async_comprehensive.c

50
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->arg = arg;
node->callback = callback; node->callback = callback;
node->ua = ua; node->ua = ua;
// Calculate expiration time in milliseconds if (timeout_tb == 0) {
struct timeval now; node->expiration_ms = 0;
get_current_time(&now); node->next = NULL;
timeval_add_tb(&now, timeout_tb); if (ua->immediate_queue_tail) ua->immediate_queue_tail->next = node;
node->expiration_ms = timeval_to_ms(&now); 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 // Add to heap
if (timeout_heap_push(ua->timeout_heap, node->expiration_ms, node) != 0) { 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"); DEBUG_ERROR(DEBUG_CATEGORY_TIMERS, "uasync_set_timeout: failed to push to heap");
memory_pool_free(ua->timeout_pool, node); memory_pool_free(ua->timeout_pool, node);
ua->timer_free_count++; // Balance the alloc counter ua->timer_free_count++; // Balance the alloc counter
return NULL; return NULL;
} }
return node; return node;
} }
// Immediate execution in next mainloop (FIFO order) // 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; 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", 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); ua, t_id, node, (unsigned long long)node->expiration_ms);
// If not found in heap, it may have already expired or been invalid // 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; return ERR_FAIL;
} }

11
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; struct dummynet_pkt* pkt = (struct dummynet_pkt*)entry->data;
ssize_t n = socket_sendto(dn->sock, pkt->data, pkt->len, ssize_t n = socket_sendto(dn->sock, pkt->data, pkt->len,
(struct sockaddr*)&dir->dest_addr, dir->dest_len); (struct sockaddr*)&dir->dest_addr, dir->dest_len);
if (n < 0) { if (n < 0) {
@ -327,7 +326,15 @@ static void dummynet_delay_callback(void* user_arg) {
if (sarg) { if (sarg) {
sarg->dn = dn; sarg->dn = dn;
sarg->dir_idx = dir_idx; 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 { } else {
dir->shaper_pending = 0; dir->shaper_pending = 0;
} }

14
tests/test_dummynet.c

@ -60,7 +60,7 @@ struct test_state {
int passed; int passed;
}; };
static inline uint64_t get_time_us(void) { static inline uint64_t now_us(void) {
struct timeval tv; struct timeval tv;
gettimeofday(&tv, NULL); gettimeofday(&tv, NULL);
return (uint64_t)tv.tv_sec * 1000000ULL + tv.tv_usec; 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)); uint16_t data_len = MIN_PKT_SIZE + (rand() % (max_data_len - MIN_PKT_SIZE + 1));
pkt->seq_num = seq; pkt->seq_num = seq;
pkt->send_time_us = get_time_us(); pkt->send_time_us = now_us();
pkt->data_len = data_len; pkt->data_len = data_len;
/* Заполняем данные */ /* Заполняем данные */
@ -95,7 +95,7 @@ static void send_one_packet(struct test_state* st, uint32_t seq) {
/* B -> dummynet (backward) */ /* B -> dummynet (backward) */
pkt->seq_num = seq + 10000; /* Отличаем пакеты от B */ 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, socket_sendto(st->sock_b, buf, sizeof(*pkt) + data_len,
(struct sockaddr*)&dest, sizeof(dest)); (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; if (n < (ssize_t)sizeof(struct test_pkt)) return;
struct test_pkt* pkt = (struct test_pkt*)buf; 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; uint64_t latency = now - pkt->send_time_us;
st->recv_a++; 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; if (n < (ssize_t)sizeof(struct test_pkt)) return;
struct test_pkt* pkt = (struct test_pkt*)buf; 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; uint64_t latency = now - pkt->send_time_us;
st->recv_b++; 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"); st.send_timer = uasync_set_timeout(st.ua, 1, &st, send_timer_callback, "test_dummynet");
/* Умное ожидание завершения: ранний выход при готовности + watchdog */ /* Умное ожидание завершения: ранний выход при готовности + watchdog */
uint64_t start_time = get_time_us(); uint64_t start_time = now_us();
uint64_t watchdog_start = 0; /* Стартует после отправки всех пакетов */ uint64_t watchdog_start = 0; /* Стартует после отправки всех пакетов */
int watchdog_active = 0; int watchdog_active = 0;
const uint64_t WATCHDOG_US = 3000000ULL; /* 3 секунды максимум после отправки */ const uint64_t WATCHDOG_US = 3000000ULL; /* 3 секунды максимум после отправки */
@ -273,7 +273,7 @@ static int run_scenario(const char* name,
while (1) { while (1) {
uasync_poll(st.ua, 1); /* 0.1ms timeout - чаще для обработки таймеров */ uasync_poll(st.ua, 1); /* 0.1ms timeout - чаще для обработки таймеров */
uint64_t now = get_time_us(); uint64_t now = now_us();
/* Активируем watchdog когда все пакеты отправлены */ /* Активируем watchdog когда все пакеты отправлены */
if (!watchdog_active && st.next_seq >= PKT_COUNT) { if (!watchdog_active && st.next_seq >= PKT_COUNT) {

52
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"); 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"); ASSERT_EQ(ctx.callback_count, 0, "Callbacks fired too early");
uasync_poll(ua, 1);
/* 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 */
}
ASSERT_EQ(ctx.callback_count, ctx.expected_count, "Immediate timeouts didn't fire correctly"); ASSERT_EQ(ctx.callback_count, ctx.expected_count, "Immediate timeouts didn't fire correctly");
/* Check that exactly 5 new immediate timeouts were recorded */ /* Check that exactly 5 new immediate timeouts were recorded */
@ -260,6 +255,48 @@ static void test_immediate_timeouts(void) {
TEST_PASS(); 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 */ /* Test 4: Memory leak detection */
static void test_memory_leak_detection(void) { static void test_memory_leak_detection(void) {
TEST_START("Memory leak detection"); TEST_START("Memory leak detection");
@ -502,6 +539,7 @@ int main(void) {
test_basic_timers(); test_basic_timers();
test_timer_cancellation_races(); test_timer_cancellation_races();
test_immediate_timeouts(); test_immediate_timeouts();
test_timeout_ordering();
test_memory_leak_detection(); test_memory_leak_detection();
test_socket_management(); test_socket_management();
test_error_handling(); test_error_handling();

Loading…
Cancel
Save