From 6a971389b6f886e97eda53abfd37ef3d773ee3f5 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 22 Apr 2026 00:04:20 +0300 Subject: [PATCH] channel select: round-robin and inflight --- src/etcp.c | 2 ++ src/etcp.h | 1 + src/etcp_connections.c | 1 + src/etcp_loadbalancer.c | 79 +++++++++++++++++++++++++++-------------- 4 files changed, 57 insertions(+), 26 deletions(-) diff --git a/src/etcp.c b/src/etcp.c index 93acb542..f8f35601 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -97,6 +97,7 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n etcp->optimal_inflight=100000; etcp->initialized=0; etcp->links_up=0; + etcp->last_rr_link=NULL; etcp->name = u_strdup(name); // Initialize log_name with local node_id (peer will be updated later when known) snprintf(etcp->log_name, sizeof(etcp->log_name), "%04X->???? [%s]", (uint16_t)instance->node_id, etcp->name); @@ -259,6 +260,7 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) { // Reset RTT history memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history)); etcp->rtt_history_idx = 0; + etcp->last_rr_link = NULL; // Clear queues (keep queue structures) clear_queue(etcp->input_queue); diff --git a/src/etcp.h b/src/etcp.h index 0c4085cf..01aa3f39 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -77,6 +77,7 @@ struct ETCP_CONN { // Links (channels) - linked list struct ETCP_LINK* links; + struct ETCP_LINK* last_rr_link; // последний линк, выбранный round-robin // Crypto and state struct secure_channel crypto_ctx; diff --git a/src/etcp_connections.c b/src/etcp_connections.c index e10669a3..81f189cd 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -747,6 +747,7 @@ void etcp_link_close(struct ETCP_LINK* link) { } pp = &(*pp)->next; } + if (link->etcp->last_rr_link == link) link->etcp->last_rr_link = NULL; remove_link(link->conn, link->ip_port_hash); diff --git a/src/etcp_loadbalancer.c b/src/etcp_loadbalancer.c index 53a314e0..299a5235 100644 --- a/src/etcp_loadbalancer.c +++ b/src/etcp_loadbalancer.c @@ -25,6 +25,7 @@ static void shaper_timer_cb(void* arg); // Shaper wait callback // Select link for transmission +// Algorithm: min inflight_bytes, round-robin among ties struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) { if (!etcp || !etcp->links) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] invalid parameters (etcp=%p, links=%p)", @@ -32,20 +33,13 @@ struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) { return NULL; } - struct ETCP_LINK* best = NULL; - uint64_t min_load_tb = UINT64_MAX; - uint64_t now_tb=get_time_tb(); + // First pass: find minimum inflight_bytes among available links + uint32_t min_inflight = UINT32_MAX; + int available_count = 0; struct ETCP_LINK* link = etcp->links; - int link_index = 0; while (link) { - link_index++; - -// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_select_link: link %d (%p) - initialized=%u, state=%u, load_tb=%llu, sub=%llu", -// link_index, link, link->initialized, link->shaper_timer==NULL?1:0, -// (unsigned long long)link->shaper_load_time_tb, (unsigned long long)link->shaper_sub_nanotime); - - if (!link->initialized || link->link_status != 1) {// link up/down + if (!link->initialized || link->link_status != 1) { link = link->next; continue; } @@ -54,28 +48,61 @@ struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) { continue; } - if (link->shaper_load_time_tb < now_tb-ALLOWED_DELTA) link->shaper_load_time_tb = now_tb-ALLOWED_DELTA; - - // Check if ready (load < now + burst allowance) - uint64_t effective_load_ns = link->shaper_load_time_tb * TIMEBASE_NS + (link->shaper_sub_nanotime / 10); // To ns - uint64_t now_ns = now_tb * TIMEBASE_NS; - int64_t delta_t=(int64_t)(now_ns - effective_load_ns);// uint64 может переполняться. чтобы обеспечить правильную цикличность. - if (link->shaper_load_time_tb < min_load_tb || (link->shaper_load_time_tb == min_load_tb && link->shaper_sub_nanotime < best->shaper_sub_nanotime)) { - min_load_tb = link->shaper_load_time_tb; - best = link; + available_count++; + if (link->inflight_bytes < min_inflight) { + min_inflight = link->inflight_bytes; } link = link->next; } - if (best) { -// DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] selected link %p (load_tb=%llu)", -// etcp->log_name, best, (unsigned long long)min_load_tb); - } else { + if (available_count == 0) { if (etcp->links) - DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no suitable link found: inf:%s, tmr:%s", etcp->log_name, etcp->links->send_blocked_inflight?"wait":"rdy", etcp->links->shaper_timer?"wait":"rdy"); + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no suitable link found: inf:%s, tmr:%s", + etcp->log_name, etcp->links->send_blocked_inflight?"wait":"rdy", etcp->links->shaper_timer?"wait":"rdy"); + return NULL; } - return best; + // Count how many links have the minimum inflight + int min_count = 0; + link = etcp->links; + while (link) { + if (link->initialized && link->link_status == 1 && + loadbalancer_link_can_send(link) && link->inflight_bytes == min_inflight) { + min_count++; + } + link = link->next; + } + + if (min_count == 1) { + // Only one link with minimum inflight - return it directly + link = etcp->links; + while (link) { + if (link->initialized && link->link_status == 1 && + loadbalancer_link_can_send(link) && link->inflight_bytes == min_inflight) { + return link; + } + link = link->next; + } + } + + // Round-robin among links with minimum inflight + // Start from the link after last_rr_link (or from the beginning if last_rr_link is NULL) + struct ETCP_LINK* start = etcp->last_rr_link ? etcp->last_rr_link->next : etcp->links; + if (!start) start = etcp->links; + + link = start; + do { + if (link->initialized && link->link_status == 1 && + loadbalancer_link_can_send(link) && link->inflight_bytes == min_inflight) { + etcp->last_rr_link = link; + return link; + } + link = link->next; + if (!link) link = etcp->links; + } while (link != start); + + // Should not reach here, but return NULL just in case + return NULL; } // New: Send dgram (select link, encrypt/send, update shaper)