diff --git a/lib/ll_queue.c b/lib/ll_queue.c index bbc878f1..339830b8 100644 --- a/lib/ll_queue.c +++ b/lib/ll_queue.c @@ -291,13 +291,15 @@ static void check_waiters(struct ll_queue* q) { struct queue_waiter_handle* h = waiter->handle; if (h) h->internal = NULL; - + + queue_threshold_callback_fn cb = waiter->callback; + void* arg = waiter->callback_arg; u_free(waiter); - - if (h) { + + 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); - h->callback(q, h->callback_arg); + cb(q, arg); } } @@ -554,29 +556,18 @@ void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes) q->threshold_max_bytes = max_bytes; } -void queue_waiter_handle_init(struct queue_waiter_handle* h, - queue_threshold_callback_fn callback, void* arg) { - if (!h) return; - h->internal = NULL; - h->callback = callback; - h->callback_arg = arg; -} +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; + + if (h->internal) queue_waiter_cancel(q, h); -int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h) { - if (!q || !h) return 1; - - if (h->internal) { - DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] queue_waiter_wait: waiter already registered", q->name); - return 1; - } - - // Проверить условие немедленно if (q->count <= q->threshold_max_packets && (q->threshold_max_bytes == 0 || q->total_bytes <= q->threshold_max_bytes)) { - h->callback(q, h->callback_arg); + callback(q, arg); return 1; } - + struct queue_waiter* waiter = u_malloc(sizeof(struct queue_waiter)); if (!waiter) { DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] queue_waiter_wait: u_malloc failed", q->name); @@ -584,11 +575,13 @@ int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h) { } waiter->next = NULL; waiter->handle = h; + waiter->callback = callback; + waiter->callback_arg = arg; h->internal = waiter; - + if (q->waiter_tail) q->waiter_tail->next = waiter; else q->waiter_head = waiter; q->waiter_tail = waiter; - + return 0; } diff --git a/lib/ll_queue.h b/lib/ll_queue.h index 61b3c805..76ff8531 100644 --- a/lib/ll_queue.h +++ b/lib/ll_queue.h @@ -99,6 +99,8 @@ typedef void (*queue_threshold_callback_fn)(struct ll_queue* q, void* arg); struct queue_waiter { struct queue_waiter* next; ///< Следующий в списке ожидающих struct queue_waiter_handle* handle; ///< Обратная ссылка на публичную структуру + queue_threshold_callback_fn callback; + void* callback_arg; }; /** @@ -106,15 +108,12 @@ struct queue_waiter { * @brief Публичная управляющая структура (встраивается в структуру вызывающей стороны). * * Вызывающая сторона: - * 1. Объявляет поле struct queue_waiter_handle в своей структуре - * 2. Вызывает queue_waiter_handle_init() один раз при инициализации + * 1. Объявляет поле struct queue_waiter_handle в своей структуре (= {0} или u_calloc) * 3. Вызывает queue_waiter_wait() когда хочет дождаться освобождения очереди * 4. Вызывает queue_waiter_cancel() при деинициализации для отмены ожидания */ struct queue_waiter_handle { struct queue_waiter* internal; ///< NULL = не ждём; иначе — внутренний узел - queue_threshold_callback_fn callback; - void* callback_arg; }; /** @@ -223,22 +222,18 @@ void queue_resume_callback(struct ll_queue* q); */ void queue_set_threshold(struct ll_queue* q, int max_packets, size_t max_bytes); -/** - * @brief Инициализирует публичную управляющую структуру waiter. - * @param h указатель на handle (встраивается в структуру вызывающей стороны) - * @param callback функция, вызываемая при достижении порога - * @param arg аргумент коллбэка - */ -void queue_waiter_handle_init(struct queue_waiter_handle* h, - queue_threshold_callback_fn callback, void* arg); - /** * @brief Регистрирует ожидание освобождения очереди до общего порога. * @param q очередь - * @param h указатель на инициализированный handle - * @return 1 — условие уже выполнено (callback вызван немедленно), 0 — зарегистрирован в списке + * @param h указатель на handle (создаётся как `struct queue_waiter_handle h = {0}` или `h.internal = NULL`) + * @param callback функция, вызываемая при достижении порога + * @param arg аргумент коллбэка + * @return 1 — условие уже выполнено (callback вызван немедленно), 0 — зарегистрирован в списке, -1 — ошибка + * + * Если handle уже ожидает (h->internal != NULL) — старый waiter отменяется и создаётся новый. */ -int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h); +int queue_waiter_wait(struct ll_queue* q, struct queue_waiter_handle* h, + queue_threshold_callback_fn callback, void* arg); /** * @brief Отменяет ожидание (удаляет из списка, если зарегистрирован). diff --git a/src/etcp_router.c b/src/etcp_router.c index 88c30842..517f06b7 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -516,11 +516,12 @@ int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id) { } void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, - struct queue_waiter_handle* h) { + struct queue_waiter_handle* h, + queue_threshold_callback_fn callback, void* arg) { if (!inst || !inst->bgp || !h) return; struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, peer_node_id); if (!conn || !conn->normalizer || !conn->normalizer->input) return; - queue_waiter_wait(conn->normalizer->input, h); + queue_waiter_wait(conn->normalizer->input, h, callback, arg); } void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, diff --git a/src/etcp_router.h b/src/etcp_router.h index 6bba766e..78cf6552 100644 --- a/src/etcp_router.h +++ b/src/etcp_router.h @@ -93,7 +93,8 @@ int etcp_router_input_q_count(struct UTUN_INSTANCE* inst, uint64_t node_id); // Backpressure: зарегистрировать/отменить waiter на normalizer->input очереди void etcp_router_waiter_register(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, - struct queue_waiter_handle* h); + struct queue_waiter_handle* h, + queue_threshold_callback_fn callback, void* arg); void etcp_router_waiter_cancel(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, struct queue_waiter_handle* h); diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index 24373767..88d6df73 100644 --- a/src/pkt_normalizer.c +++ b/src/pkt_normalizer.c @@ -69,7 +69,7 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) { pn->recvpart = NULL; pn->flush_timer = NULL; - queue_waiter_handle_init(&pn->input_waiter_handle, etcp_input_ready_cb, pn); + pn->input_waiter_handle.internal = NULL; return pn; } @@ -215,7 +215,7 @@ static void packer_cb(struct ll_queue* q, void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; if (!pn) return; DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "input_q->pn: waiting etcp input threshold"); - queue_waiter_wait(pn->etcp->input_queue, &pn->input_waiter_handle); + queue_waiter_wait(pn->etcp->input_queue, &pn->input_waiter_handle, etcp_input_ready_cb, pn); } // Helper to send block to ETCP as ETCP_FRAGMENT diff --git a/src/proxy/tcp_proxy_client.c b/src/proxy/tcp_proxy_client.c index 6dac168f..ae3e8891 100644 --- a/src/proxy/tcp_proxy_client.c +++ b/src/proxy/tcp_proxy_client.c @@ -227,7 +227,6 @@ static err_t tcp_proxy_client_accept_cb(void *arg, struct tcp_pcb *newpcb, err_t tcp_nagle_disable(newpcb); pc->to_lwip = queue_new(p->ua, 0, 0, 0, "to_lwip"); - queue_waiter_handle_init(&pc->tx_waiter, tcp_proxy_client_tx_waiter_cb, pc); if (tcp_proxy_client_send_connect(pc) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy client: send_connect failed sid=%08x to %d.%d.%d.%d:%d", pc->stream_id, pc->dest_ip[0], pc->dest_ip[1], pc->dest_ip[2], pc->dest_ip[3], ntohs(pc->dest_port)); u_free(pc); return LERR_MEM; } @@ -285,7 +284,8 @@ static err_t tcp_proxy_client_recv_cb(void *arg, struct tcp_pcb *pcb, struct pbu if (ret == 0) { tcp_recved(pcb, len); u_free(data); } else { pc->tx_buf = data; pc->tx_len = len; - etcp_router_waiter_register(pc->proxy->inst, pc->proxy->via_node_id, &pc->tx_waiter); + etcp_router_waiter_register(pc->proxy->inst, pc->proxy->via_node_id, &pc->tx_waiter, + tcp_proxy_client_tx_waiter_cb, pc); } } else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "PROXY recv malloc(%u) failed sid=%08x", len, pc->stream_id); pbuf_free(p); diff --git a/src/proxy/tcp_proxy_server.c b/src/proxy/tcp_proxy_server.c index a9738b3b..91696825 100644 --- a/src/proxy/tcp_proxy_server.c +++ b/src/proxy/tcp_proxy_server.c @@ -199,7 +199,8 @@ static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) { if (rc->pause_buf) { memcpy(rc->pause_buf, buf, n); rc->pause_len = n; } else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE malloc(%zd) failed sid=%08x", n, rc->stream_id); if (rc->read_id) { uasync_remove_socket_t(rc->ua, rc->sock); rc->read_id = NULL; } - etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter); + etcp_router_waiter_register(inst, rc->peer_node_id, &rc->pause_waiter, + tcp_proxy_server_pause_waiter_cb, rc); } } else if (n == 0) { DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:EOF fd=%d sid=%08x total=%d cli_closed=%d sock_closed=%d", @@ -298,7 +299,6 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry* rc->ctx = ctx; rc->stream_id = stream_id; rc->peer_node_id = src_node_id; memcpy(rc->dest_ip, dest_ip, 4); rc->dest_port = dest_port; rc->ua = inst->ua; rc->sock = SOCKET_INVALID; - queue_waiter_handle_init(&rc->pause_waiter, tcp_proxy_server_pause_waiter_cb, rc); rc->sock = socket(AF_INET, SOCK_STREAM, 0); if (rc->sock == SOCKET_INVALID) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "TCP proxy server: socket() failed"); u_free(rc); queue_dgram_free(entry); queue_entry_free(entry); return -1; } diff --git a/tests/test_intensive_memory_pool.c b/tests/test_intensive_memory_pool.c index e69f7fce..9d917b9c 100644 --- a/tests/test_intensive_memory_pool.c +++ b/tests/test_intensive_memory_pool.c @@ -37,12 +37,9 @@ static double test_without_pools(int iterations) { } // Создать много waiters (общий порог 0, count=10 => все в очередь) - struct queue_waiter_handle handles[32]; - for (int i = 0; i < 32; i++) { - queue_waiter_handle_init(&handles[i], intensive_waiter_callback, NULL); - queue_waiter_wait(queue, &handles[i]); - } - + struct queue_waiter_handle handles[32] = {0}; + for (int i = 0; i < 32; i++) queue_waiter_wait(queue, &handles[i], intensive_waiter_callback, NULL); + // Удалить записи (триггер waiters — по одному за get) for (int i = 0; i < 10; i++) { void* retrieved = queue_data_get(queue); @@ -50,9 +47,7 @@ static double test_without_pools(int iterations) { } // Отменить оставшиеся waiters - for (int i = 0; i < 32; i++) { - queue_waiter_cancel(queue, &handles[i]); - } + for (int i = 0; i < 32; i++) queue_waiter_cancel(queue, &handles[i]); } queue_free(queue); @@ -85,11 +80,8 @@ static double test_with_pools(int iterations) { } // Создать много waiters - struct queue_waiter_handle handles[32]; - for (int i = 0; i < 32; i++) { - queue_waiter_handle_init(&handles[i], intensive_waiter_callback, NULL); - queue_waiter_wait(queue, &handles[i]); - } + struct queue_waiter_handle handles[32] = {0}; + for (int i = 0; i < 32; i++) queue_waiter_wait(queue, &handles[i], intensive_waiter_callback, NULL); // Удалить записи (триггер waiters) for (int i = 0; i < 10; i++) { @@ -98,9 +90,7 @@ static double test_with_pools(int iterations) { } // Отменить оставшиеся waiters - for (int i = 0; i < 32; i++) { - queue_waiter_cancel(queue, &handles[i]); - } + for (int i = 0; i < 32; i++) queue_waiter_cancel(queue, &handles[i]); } queue_free(queue); diff --git a/tests/test_ll_queue.c b/tests/test_ll_queue.c index 3dc38d0e..af8c6718 100644 --- a/tests/test_ll_queue.c +++ b/tests/test_ll_queue.c @@ -177,9 +177,8 @@ static void test_waiter(void) { queue_set_threshold(q, 2, 0); int called = 0; - struct queue_waiter_handle h; - queue_waiter_handle_init(&h, waiter_cb, &called); - int ret = queue_waiter_wait(q, &h); + struct queue_waiter_handle h = {0}; + int ret = queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(ret == 1 && called == 1, "immediate when condition met"); // empty queue for (int i = 0; i < 5; i++) { @@ -188,7 +187,7 @@ static void test_waiter(void) { } called = 0; - ret = queue_waiter_wait(q, &h); + ret = queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(ret == 0 && called == 0, ""); for (int i = 0; i < 3; i++) queue_entry_free((struct ll_entry*)queue_data_get(q)); @@ -214,11 +213,8 @@ static void test_waiter_multiple(void) { waiter_order_idx = 0; int ids[4] = {10, 20, 30, 40}; - struct queue_waiter_handle handles[4]; - for (int i = 0; i < 4; i++) { - queue_waiter_handle_init(&handles[i], waiter_cb_order, &ids[i]); - ASSERT(queue_waiter_wait(q, &handles[i]) == 0, "should queue"); - } + struct queue_waiter_handle handles[4] = {0}; + for (int i = 0; i < 4; i++) ASSERT(queue_waiter_wait(q, &handles[i], waiter_cb_order, &ids[i]) == 0, "should queue"); for (int i = 0; i < 5; i++) queue_entry_free((struct ll_entry*)queue_data_get(q)); ASSERT_EQ(waiter_order_idx, 3, "3 waiters fired"); @@ -242,11 +238,8 @@ static void test_waiter_cancel(void) { waiter_order_idx = 0; int ids[3] = {1, 2, 3}; - struct queue_waiter_handle handles[3]; - for (int i = 0; i < 3; i++) { - queue_waiter_handle_init(&handles[i], waiter_cb_order, &ids[i]); - queue_waiter_wait(q, &handles[i]); // count=10 > 3, all queued - } + struct queue_waiter_handle handles[3] = {0}; + for (int i = 0; i < 3; i++) queue_waiter_wait(q, &handles[i], waiter_cb_order, &ids[i]); // count=10 > 3, all queued queue_waiter_cancel(q, &handles[1]); // cancel middle (ID 2) @@ -270,9 +263,8 @@ static void test_waiter_queue_free(void) { } int called = 0; - struct queue_waiter_handle h; - queue_waiter_handle_init(&h, waiter_cb, &called); - queue_waiter_wait(q, &h); + struct queue_waiter_handle h = {0}; + queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(h.internal != NULL, "waiter is queued"); queue_free(q); @@ -293,9 +285,8 @@ static void test_waiter_threshold(void) { } int called = 0; - struct queue_waiter_handle h; - queue_waiter_handle_init(&h, waiter_cb, &called); - int ret = queue_waiter_wait(q, &h); + struct queue_waiter_handle h = {0}; + int ret = queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(ret == 0 && called == 0, "queued (10 > 0)"); queue_set_threshold(q, 10, 0); // raise threshold @@ -313,10 +304,9 @@ static void test_waiter_reuse(void) { struct ll_queue *q = queue_new(ua, 0, 0, 0, "wr"); int called = 0; - struct queue_waiter_handle h; - queue_waiter_handle_init(&h, waiter_cb, &called); + struct queue_waiter_handle h = {0}; - int ret = queue_waiter_wait(q, &h); + int ret = queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(ret == 1 && called == 1, "immediate (empty queue)"); for (int i = 0; i < 3; i++) { @@ -325,7 +315,7 @@ static void test_waiter_reuse(void) { } called = 0; - ret = queue_waiter_wait(q, &h); + ret = queue_waiter_wait(q, &h, waiter_cb, &called); ASSERT(ret == 0 && called == 0, "queued"); for (int i = 0; i < 3; i++) queue_entry_free((struct ll_entry*)queue_data_get(q)); diff --git a/tests/test_memory_pool_and_config.c b/tests/test_memory_pool_and_config.c index 64143e81..ce066d43 100644 --- a/tests/test_memory_pool_and_config.c +++ b/tests/test_memory_pool_and_config.c @@ -45,12 +45,11 @@ int main() { return 1; } queue_set_threshold(queue, 2, 0); - + // Test multiple waiters with new linked-list backpressure - struct queue_waiter_handle handles[10]; + struct queue_waiter_handle handles[10] = {0}; for (int i = 0; i < 10; i++) { - queue_waiter_handle_init(&handles[i], test_waiter_callback, NULL); - queue_waiter_wait(queue, &handles[i]); // queue empty => immediate callback + queue_waiter_wait(queue, &handles[i], test_waiter_callback, NULL); // queue empty => immediate callback } // Add some entries — waiters already fired, none queued