@ -391,6 +391,11 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) {
clear_queue ( etcp - > recv_q ) ;
clear_queue ( etcp - > recv_q ) ;
clear_queue ( etcp - > ack_q ) ;
clear_queue ( etcp - > ack_q ) ;
// clear_queue leaves callback_suspended=1, resume to prevent deadlock
queue_resume_callback ( etcp - > input_queue ) ;
queue_resume_callback ( etcp - > input_send_q ) ;
queue_resume_callback ( etcp - > input_wait_ack ) ;
// В etcp_conn_reset(), после очистки очередей добавьте:
// В etcp_conn_reset(), после очистки очередей добавьте:
struct ETCP_LINK * l = etcp - > links ;
struct ETCP_LINK * l = etcp - > links ;
while ( l ) {
while ( l ) {
@ -612,6 +617,16 @@ static void input_queue_try_resume(struct ETCP_CONN* etcp) {// при ACK
}
}
}
}
// Called from link-level etcp_link_update_inflight_lim() after inflight_lim_bytes changes.
// Recalculates connection-level optimal_inflight and resumes input_queue if room opened up.
void etcp_conn_on_inflight_lim_changed ( struct ETCP_CONN * etcp ) {
if ( ! etcp ) return ;
uint32_t sum = 0 ;
for ( struct ETCP_LINK * tl = etcp - > links ; tl ; tl = tl - > next ) sum + = tl - > inflight_lim_bytes ;
etcp - > optimal_inflight = sum ;
input_queue_try_resume ( etcp ) ;
}
void etcp_stats ( struct ETCP_CONN * etcp ) {
void etcp_stats ( struct ETCP_CONN * etcp ) {
if ( ! etcp ) return ;
if ( ! etcp ) return ;
@ -924,7 +939,8 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
bbr_note_loss ( inf_pkt - > last_link - > bbr ) ;
bbr_note_loss ( inf_pkt - > last_link - > bbr ) ;
inf_pkt - > last_link - > inflight_bytes - = inf_pkt - > ll . len ;
inf_pkt - > last_link - > inflight_bytes - = inf_pkt - > ll . len ;
inf_pkt - > last_link - > inflight_packets - - ;
inf_pkt - > last_link - > inflight_packets - - ;
if ( link - > send_blocked_inflight & & link - > inflight_bytes < link - > inflight_lim_bytes ) loadbalancer_link_ready ( link ) ;
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)
// Always add to the CURRENT link (first send or retransmission)
@ -1239,13 +1255,9 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
uint32_t new_pacing = link - > bbr_pacing_rate ;
uint32_t new_pacing = link - > bbr_pacing_rate ;
bbr_main ( link - > bbr , & rs , & new_cwnd , & new_pacing , link - > mtu , link - > inflight_bytes ,
bbr_main ( link - > bbr , & rs , & new_cwnd , & new_pacing , link - > mtu , link - > inflight_bytes ,
( link - > inflight_bytes > = link - > inflight_lim_bytes ) ) ;
( link - > inflight_bytes > = link - > inflight_lim_bytes ) ) ;
etcp_link_update_inflight_lim ( link , new_cwnd ) ;
link - > bbr_pacing_rate = new_pacing ;
link - > bbr_pacing_rate = new_pacing ;
link - > bandwidth = ( uint32_t ) ( ( uint64_t ) new_pacing * 8 / 1000 ) ;
link - > bandwidth = ( uint32_t ) ( ( uint64_t ) new_pacing * 8 / 1000 ) ;
if ( link - > send_blocked_inflight & & link - > inflight_bytes < new_cwnd ) {
link - > send_blocked_inflight = 0 ;
loadbalancer_link_ready ( link ) ;
}
link - > bbr_loss_since_ack = 0 ;
link - > bbr_loss_since_ack = 0 ;
}
}
@ -1277,6 +1289,10 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
void etcp_conn_input ( struct ETCP_DGRAM * pkt ) {
void etcp_conn_input ( struct ETCP_DGRAM * pkt ) {
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " " ) ;
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " " ) ;
if ( ! pkt ) return ;
if ( ! pkt ) return ;
if ( memory_pool_is_freed ( pkt - > link - > etcp - > instance - > pkt_pool , pkt ) ) {
DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING " , ( void * ) pkt ) ;
volatile int _halt = 1 ; while ( _halt ) { }
}
if ( ! pkt - > data_len ) {
if ( ! pkt - > data_len ) {
memory_pool_free ( pkt - > link - > etcp - > instance - > pkt_pool , pkt ) ;
memory_pool_free ( pkt - > link - > etcp - > instance - > pkt_pool , pkt ) ;
return ;
return ;
@ -1387,24 +1403,22 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
break ;
break ;
}
}
}
}
if ( queue_find_data_by_index ( etcp - > ack_q , & seq ) = = NULL ) {
struct ACK_PACKET * p = ( struct ACK_PACKET * ) queue_entry_new_from_pool ( etcp - > instance - > ack_pool ) ;
struct ACK_PACKET * p = ( struct ACK_PACKET * ) queue_entry_new_from_pool ( etcp - > instance - > ack_pool ) ;
if ( ! p ) {
if ( ! p ) {
DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " [%s] failed to allocate ACK_PACKET " , etcp - > log_name ) ;
DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " [%s] failed to allocate ACK_PACKET " , etcp - > log_name ) ;
len = 0 ;
len = 0 ;
break ;
break ;
}
}
p - > seq = seq ;
p - > seq = seq ;
p - > pkt_timestamp = pkt - > timestamp ;
p - > pkt_timestamp = pkt - > timestamp ;
p - > recv_timestamp = get_current_timestamp ( ) ;
p - > recv_timestamp = get_current_timestamp ( ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] RX add to ack_q seq=%d " , etcp - > log_name , seq ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] RX add to ack_q seq=%d " , etcp - > log_name , seq ) ;
queue_data_put_with_index ( etcp - > ack_q , ( struct ll_entry * ) p ) ;
queue_data_put_with_index ( etcp - > ack_q , ( struct ll_entry * ) p ) ;
if ( etcp - > ack_resp_timer = = NULL ) {
if ( etcp - > ack_resp_timer = = NULL ) {
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " [%s] set ack_timer for delayed ACK send " , etcp - > log_name ) ;
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " [%s] set ack_timer for delayed ACK send " , etcp - > log_name ) ;
etcp - > ack_resp_timer = uasync_set_timeout ( etcp - > instance - > ua , ACK_DELAY_TB , etcp , ack_response_timer_cb , " etcp_ack_resp " ) ;
etcp - > ack_resp_timer = uasync_set_timeout ( etcp - > instance - > ua , ACK_DELAY_TB , etcp , ack_response_timer_cb , " etcp_ack_resp " ) ;
}
}
} else DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] RX ack dedup: seq=%d already in ack_q " , etcp - > log_name , seq ) ;
if ( ( ( int32_t ) ( etcp - > last_delivered_id - seq ) < 0 ) & & ( queue_find_data_by_index ( etcp - > recv_q , & seq ) = = NULL ) ) { // проверяем есть ли пакет с этим seq
if ( ( ( int32_t ) ( etcp - > last_delivered_id - seq ) < 0 ) & & ( queue_find_data_by_index ( etcp - > recv_q , & seq ) = = NULL ) ) { // проверяем есть ли пакет с этим seq
uint32_t pkt_len = len - 5 ;
uint32_t pkt_len = len - 5 ;
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " [%s] adding packet seq=%u to recv_q (last_delivered_id=%u) " , etcp - > log_name , seq , etcp - > last_delivered_id ) ;
DEBUG_TRACE ( DEBUG_CATEGORY_ETCP , " [%s] adding packet seq=%u to recv_q (last_delivered_id=%u) " , etcp - > log_name , seq , etcp - > last_delivered_id ) ;
@ -1546,6 +1560,12 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
}
}
if ( memory_pool_is_freed ( pkt - > link - > etcp - > instance - > pkt_pool , pkt ) ) {
DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " etcp_conn_input: pkt=%p ALREADY FREED in pkt_pool — HALTING " , ( void * ) pkt ) ;
volatile int _halt = 1 ; while ( _halt ) { }
}
memory_pool_free ( etcp - > instance - > pkt_pool , pkt ) ; // Free the incoming dgram
memory_pool_free ( etcp - > instance - > pkt_pool , pkt ) ; // Free the incoming dgram
}
}