From 914f252b1283aeed12881a5cabd5664cc207d938 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Mon, 29 Jun 2026 20:04:20 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20remove=20send=5Fblocked=5Finflight=20?= =?UTF-8?q?=E2=80=94=20per-link=20inflight=20gate=20was=20blocking=20input?= =?UTF-8?q?=5Fsend=5Fq=20and=20ACK=20send,=20causing=20bidirectional=20dea?= =?UTF-8?q?dlock?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/etcp.c | 4 +--- src/etcp_connections.c | 10 +--------- src/etcp_connections.h | 3 --- src/etcp_dump.c | 3 +-- src/etcp_loadbalancer.c | 24 ++++++------------------ tests/test_etcp_reinit_inflight.c | 7 ++++--- 6 files changed, 13 insertions(+), 38 deletions(-) diff --git a/src/etcp.c b/src/etcp.c index 6ef6cfd0..ac83132e 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -864,7 +864,7 @@ static void etcp_conn_process_send_queue(struct ETCP_CONN* etcp) {// вызыв char l_status[256]={0}; struct ETCP_LINK* link = etcp->links; while (link) { - snprintf (l_status+strlen(l_status), 256-strlen(l_status), "L%d%d%s,%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy", link->send_blocked_inflight?"inf_block":"rdy"); + snprintf (l_status+strlen(l_status), 256-strlen(l_status), "L%d%d%s ", link->recv_keepalive, link->remote_keepalive, (link->shaper_timer)?"wait":"rdy"); link = link->next; } DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TX state: %d (link not ready, skip send) %s", etcp->log_name, etcp->tx_state, l_status); @@ -934,8 +934,6 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { bbr_note_loss(inf_pkt->last_link->bbr); inf_pkt->last_link->inflight_bytes -= inf_pkt->ll.len; inf_pkt->last_link->inflight_packets--; - if (inf_pkt->last_link->send_blocked_inflight && inf_pkt->last_link->inflight_bytes < inf_pkt->last_link->inflight_lim_bytes) - loadbalancer_link_ready(inf_pkt->last_link); } // Always add to the CURRENT link (first send or retransmission) diff --git a/src/etcp_connections.c b/src/etcp_connections.c index c50aeee1..a70cac94 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -869,16 +869,8 @@ void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) { if (new_lim < INFLIGHT_LIM_MIN) new_lim = INFLIGHT_LIM_MIN; if (new_lim > link->etcp->max_inflight) new_lim = link->etcp->max_inflight; - uint32_t old = link->inflight_lim_bytes; link->inflight_lim_bytes = new_lim; - - if (link->inflight_bytes < new_lim && link->send_blocked_inflight) {// разблокируем линк - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] unblocking link %p (inflight_lim %u -> %u)", link->etcp->log_name, link, old, new_lim); - etcp_conn_on_inflight_lim_changed(link->etcp);// сперва обновим inflight в etcp - link->send_blocked_inflight = 0; - if (!link->shaper_timer) loadbalancer_link_ready(link);// потом вызовем ready -> разблок - } else etcp_conn_on_inflight_lim_changed(link->etcp);// - + etcp_conn_on_inflight_lim_changed(link->etcp); } void etcp_link_close(struct ETCP_LINK* link) { diff --git a/src/etcp_connections.h b/src/etcp_connections.h index 638d04e7..d0c8e710 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -171,9 +171,6 @@ struct ETCP_LINK { uint32_t inflight_packets; uint32_t inflight_lim_bytes; - /* Биты блокировки отправки по причине (для корректного resume) */ - unsigned send_blocked_inflight : 1; - // statistics size_t encrypt_errors; size_t decrypt_errors; diff --git a/src/etcp_dump.c b/src/etcp_dump.c index 0de1046c..9cb372a6 100644 --- a/src/etcp_dump.c +++ b/src/etcp_dump.c @@ -161,9 +161,8 @@ void etcp_dump_conn_state(struct ETCP_CONN* conn) { link->keepalive_sent_count > 0 ? "has_sent" : "idle"); /* inflight (BBR) */ - DUMP_LINE("BBR: bytes=%u pkts=%u lim=%u blocked=%d mode=%d cycle=%d pacing=%u", + DUMP_LINE("BBR: bytes=%u pkts=%u lim=%u mode=%d cycle=%d pacing=%u", link->inflight_bytes, link->inflight_packets, link->inflight_lim_bytes, - link->send_blocked_inflight, link->bbr ? link->bbr->mode : -1, link->bbr ? link->bbr->cycle_idx : -1, link->bbr_pacing_rate); diff --git a/src/etcp_loadbalancer.c b/src/etcp_loadbalancer.c index 7bae431e..d1803925 100644 --- a/src/etcp_loadbalancer.c +++ b/src/etcp_loadbalancer.c @@ -57,8 +57,8 @@ struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) { 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: tmr:%s", + etcp->log_name, etcp->links->shaper_timer?"wait":"rdy"); return NULL; } @@ -181,20 +181,8 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) { int loadbalancer_link_can_send(struct ETCP_LINK* link) { if (!link) return 0; if (link->burst_active) return 1; // burst bypasses all limits - int can = 1; - if (link->inflight_bytes >= link->inflight_lim_bytes) { - link->send_blocked_inflight = 1; - can = 0; - } else { - link->send_blocked_inflight = 0; - } - if (link->shaper_timer) { -// link->send_blocked_bandwidth = 1; - can = 0; - } else { -// link->send_blocked_bandwidth = 0; - } - return can; + if (link->shaper_timer) return 0; + return 1; } void loadbalancer_link_ready(struct ETCP_LINK* link) { @@ -202,8 +190,8 @@ void loadbalancer_link_ready(struct ETCP_LINK* link) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "invalid link (%p)", link); return; } - if (loadbalancer_link_can_send(link) == 0) { - DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "link still blocked"); + if (link->shaper_timer) { + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "link still blocked by shaper"); return; } DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "link=%p now ready, notifying ETCP_CONN", link); diff --git a/tests/test_etcp_reinit_inflight.c b/tests/test_etcp_reinit_inflight.c index c86a9922..5b3a1c31 100644 --- a/tests/test_etcp_reinit_inflight.c +++ b/tests/test_etcp_reinit_inflight.c @@ -144,7 +144,7 @@ static int buffers_full(struct test_ctx* ctx) { struct ETCP_CONN* conn = ctx->sender->connections; if (!conn) return 0; struct ETCP_LINK* link = conn->links; - return link && link->send_blocked_inflight + return link && link->inflight_bytes >= link->inflight_lim_bytes && conn->input_wait_ack->count > 5 && ctx->total_packets_sent >= TOTAL_PACKETS; } @@ -190,8 +190,9 @@ static void monitor(void* arg) { break; case 2: // fill buffers if (buffers_full(ctx)) { - printf(" buffers full: inflight_blocked=%d wait_ack=%d sent=%d\n", - ctx->sender->connections->links->send_blocked_inflight, + printf(" buffers full: inflight=%u/%u wait_ack=%d sent=%d\n", + ctx->sender->connections->links->inflight_bytes, + ctx->sender->connections->links->inflight_lim_bytes, ctx->sender->connections->input_wait_ack->count, ctx->total_packets_sent); ctx->phase = 3;