Browse Source

feat: queue_set_waiter_defer — immediate or call_soon callback mode

queue_waiter_handle stores defer_q/cb/arg, no extra malloc.
queue_set_waiter_defer(q, 1) enables call_soon dispatch for waiter callbacks.
bbr_integration uses it instead of manual send_waiter_cb wrapper.
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
2a259191b9
  1. 25
      lib/ll_queue.c
  2. 6
      lib/ll_queue.h
  3. 11
      tests/bbr_integration/test_bbr_integration.c

25
lib/ll_queue.c

@ -60,6 +60,7 @@ struct ll_queue* queue_new(struct UASYNC* ua, size_t hash_size, uint16_t index_o
q->size_limit = -1; // Без ограничения по умолчанию
q->threshold_max_packets = 0; // По умолчанию: ждать пустой очереди
q->threshold_max_bytes = 0; // По умолчанию: не проверять байты
q->waiter_defer = 0;
q->hash_size = hash_size;
q->index_offset = index_offset;
q->index_size = index_size;
@ -277,6 +278,21 @@ static uint32_t make_hash(const void* data, uint16_t len) {// алгоритм F
return hash;
}
static void waiter_defer_cb(void* arg) {
struct queue_waiter_handle* h = (struct queue_waiter_handle*)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,
struct queue_waiter_handle* h) {
if (q->waiter_defer) {
h->defer_q = q; h->defer_cb = cb; h->defer_arg = arg;
uasync_call_soon(q->ua, h, waiter_defer_cb);
} else {
cb(q, arg);
}
}
// Проверить и запустить первый ожидающий коллбэк (FIFO / round-robin)
static void check_waiters(struct ll_queue* q) {
if (!q || !q->waiter_head) return;
@ -299,7 +315,7 @@ static void check_waiters(struct ll_queue* q) {
if (cb) {
DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "check_waiters: waking head waiter, count=%d<=%d, bytes=%zu<=%zu",
q->count, q->threshold_max_packets, q->total_bytes, q->threshold_max_bytes);
cb(q, arg);
queue_waiter_call(q, cb, arg, h);
}
}
@ -592,6 +608,11 @@ void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes)
q->threshold_max_bytes = max_bytes;
}
void queue_set_waiter_defer(struct ll_queue* q, int enable) {
if (!q) return;
q->waiter_defer = enable;
}
int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h,
queue_threshold_callback_fn callback, void* arg) {
if (!q || !h) return -1;
@ -600,7 +621,7 @@ int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h,
if (q->count <= q->threshold_max_packets
&& (q->threshold_max_bytes == 0 || q->total_bytes <= q->threshold_max_bytes)) {
callback(q, arg);
queue_waiter_call(q, callback, arg, h);
return 1;
}

6
lib/ll_queue.h

@ -114,6 +114,9 @@ struct queue_waiter {
*/
struct queue_waiter_handle {
struct queue_waiter* internal; ///< NULL = не ждём; иначе — внутренний узел
struct ll_queue* defer_q;
queue_threshold_callback_fn defer_cb;
void* defer_arg;
};
/**
@ -139,6 +142,7 @@ struct ll_queue {
struct queue_waiter* waiter_tail; // Хвост списка ожидающих
int threshold_max_packets; // Общий порог: макс. кол-во элементов (по умолчанию 0)
size_t threshold_max_bytes; // Общий порог: макс. объём данных (0 = не проверять)
int waiter_defer; // 0=immediate callback, 1=defer via uasync_call_soon
struct ll_entry** hash_table; // Хеш-таблица для поиска по id (если hash_size > 0)
size_t hash_size; // Размер хеш-таблицы
@ -222,6 +226,8 @@ void queue_resume_callback(struct ll_queue* q);
*/
void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes);
void queue_set_waiter_defer(struct ll_queue* q, int enable);
/**
* @brief Регистрирует ожидание освобождения очереди до общего порога.
* @param q очередь

11
tests/bbr_integration/test_bbr_integration.c

@ -177,10 +177,9 @@ static void send_timer_cb(void* arg) {
send_burst((struct test_ctx*)arg);
}
static void send_waiter_cb(struct ll_queue* q, void* arg) {
static void send_wtr_cb(struct ll_queue* q, void* arg) {
(void)q;
struct test_ctx* ctx = (struct test_ctx*)arg;
uasync_call_soon(ctx->ua, ctx, send_timer_cb);
send_burst((struct test_ctx*)arg);
}
static void send_burst(struct test_ctx* ctx) {
@ -191,7 +190,7 @@ static void send_burst(struct test_ctx* ctx) {
if (conn->normalizer && conn->normalizer->input) {
if (queue_entry_count(conn->normalizer->input) >= SEND_QUEUE_THRESHOLD) {
queue_waiter_wait(conn->normalizer->input, &ctx->waiter,
send_waiter_cb, ctx);
send_wtr_cb, ctx);
return;
}
}
@ -391,8 +390,10 @@ int main(void) {
t0 = now_us();
while ((now_us() - t0) < 500000ULL) uasync_poll(ctx.ua, 1);
if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input)
if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input) {
queue_set_threshold(ctx.sender->connections->normalizer->input, 0, 0);
queue_set_waiter_defer(ctx.sender->connections->normalizer->input, 1);
}
/* Start test */
printf("\n=== Starting traffic (%d seconds) ===\n\n", TEST_DURATION_MS / 1000);

Loading…
Cancel
Save