@ -33,7 +33,6 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg);
static void pause_resume_cb ( struct ll_queue * q , void * arg ) ;
static void pause_resume_cb ( struct ll_queue * q , void * arg ) ;
static void close_retry_cb ( void * arg ) ;
static void close_retry_cb ( void * arg ) ;
static void diag_timer_cb ( void * arg ) ;
static void diag_timer_cb ( void * arg ) ;
static void retry_timer_cb ( void * arg ) ;
static int send_msg ( struct UTUN_INSTANCE * inst , uint64_t dst , uint8_t subcmd ,
static int send_msg ( struct UTUN_INSTANCE * inst , uint64_t dst , uint8_t subcmd ,
uint32_t sid , const uint8_t * data , size_t len , int force ) ;
uint32_t sid , const uint8_t * data , size_t len , int force ) ;
static void send_close ( struct tcp_proxy_server_conn * rc ) ;
static void send_close ( struct tcp_proxy_server_conn * rc ) ;
@ -218,14 +217,14 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
int ret = send_msg ( inst , rc - > peer_node_id , TCP_PROXY_SUBCMD_DATA , rc - > stream_id , rc - > tx_buf , rc - > tx_len , 0 ) ;
int ret = send_msg ( inst , rc - > peer_node_id , TCP_PROXY_SUBCMD_DATA , rc - > stream_id , rc - > tx_buf , rc - > tx_len , 0 ) ;
if ( ret = = 0 ) {
if ( ret = = 0 ) {
u_free ( rc - > tx_buf ) ; rc - > tx_buf = NULL ; rc - > tx_len = 0 ;
u_free ( rc - > tx_buf ) ; rc - > tx_buf = NULL ; rc - > tx_len = 0 ;
if ( rc - > retry_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > retry_timer ) ; rc - > retry_timer = NULL ; }
if ( rc - > dst_fin_deferred & & ! q - > head ) {
if ( rc - > dst_fin_deferred & & ! q - > head ) {
rc - > dst_fin_deferred = 0 ;
rc - > dst_fin_deferred = 0 ;
send_fin ( rc ) ;
send_fin ( rc ) ;
return ;
return ;
}
}
} else {
} else {
if ( ! rc - > retry_timer ) rc - > retry_timer = uasync_set_timeout ( rc - > ua , 5000 , rc , retry_timer_cb , " tps_retry " ) ;
etcp_router_on_send_ready ( inst , rc - > peer_node_id , ETCP_ID_TCP_PROXY_CLIENT ,
& rc - > pause_waiter , pause_resume_cb , rc ) ;
return ;
return ;
}
}
}
}
@ -262,7 +261,8 @@ static void read_queue_drain_cb(struct ll_queue* q, void* arg) {
memory_pool_free ( rc - > tc - > data_pool , e - > dgram ) ; queue_entry_free ( e ) ;
memory_pool_free ( rc - > tc - > data_pool , e - > dgram ) ; queue_entry_free ( e ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_SOCKET , " SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry " ,
DEBUG_DEBUG ( DEBUG_CATEGORY_SOCKET , " SOCK:BACKPRESSURE fd=%d sid=%08x — buffered %u for retry " ,
( int ) rc - > tc - > sock , rc - > stream_id , rc - > tx_len ) ;
( int ) rc - > tc - > sock , rc - > stream_id , rc - > tx_len ) ;
if ( ! rc - > retry_timer ) rc - > retry_timer = uasync_set_timeout ( rc - > ua , 5000 , rc , retry_timer_cb , " tps_retry " ) ;
etcp_router_on_send_ready ( inst , rc - > peer_node_id , ETCP_ID_TCP_PROXY_CLIENT ,
& rc - > pause_waiter , pause_resume_cb , rc ) ;
}
}
}
}
@ -273,25 +273,6 @@ static void pause_resume_cb(struct ll_queue* q, void* arg) {
queue_resume_callback ( rc - > tc - > read_queue ) ;
queue_resume_callback ( rc - > tc - > read_queue ) ;
}
}
static void retry_timer_cb ( void * arg ) {
struct tcp_proxy_server_conn * rc = ( struct tcp_proxy_server_conn * ) arg ;
rc - > retry_timer = NULL ;
if ( ! rc - > tx_buf | | ! rc - > tc | | rc - > tc - > sock = = SOCKET_INVALID | | rc - > cli_closed ) return ;
struct UTUN_INSTANCE * inst = rc - > ctx ? rc - > ctx - > inst : NULL ;
int ret = send_msg ( inst , rc - > peer_node_id , TCP_PROXY_SUBCMD_DATA , rc - > stream_id , rc - > tx_buf , rc - > tx_len , 1 ) ;
if ( ret = = 0 ) {
u_free ( rc - > tx_buf ) ; rc - > tx_buf = NULL ; rc - > tx_len = 0 ;
if ( rc - > dst_fin_deferred & & rc - > tc & & rc - > tc - > read_queue & & ! rc - > tc - > read_queue - > head ) {
rc - > dst_fin_deferred = 0 ;
send_fin ( rc ) ;
return ;
}
queue_resume_callback ( rc - > tc - > read_queue ) ;
} else {
rc - > retry_timer = uasync_set_timeout ( rc - > ua , 5000 , rc , retry_timer_cb , " tps_retry " ) ;
}
}
// ====================================================================
// ====================================================================
// Управление жизненным циклом коннекта
// Управление жизненным циклом коннекта
// ====================================================================
// ====================================================================
@ -323,8 +304,8 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
struct tcp_proxy_server_conn * * prev = & rc - > ctx - > conns ;
struct tcp_proxy_server_conn * * prev = & rc - > ctx - > conns ;
while ( * prev ) { if ( * prev = = rc ) { * prev = rc - > next ; rc - > ctx - > conn_count - - ; break ; } prev = & ( * prev ) - > next ; }
while ( * prev ) { if ( * prev = = rc ) { * prev = rc - > next ; rc - > ctx - > conn_count - - ; break ; } prev = & ( * prev ) - > next ; }
}
}
if ( rc - > ctx & & rc - > ctx - > inst ) etcp_router_cancel_send_ready ( rc - > ctx - > inst , rc - > peer_node_id , ETCP_ID_TCP_PROXY_CLIENT , & rc - > pause_waiter ) ;
if ( rc - > tx_buf ) { u_free ( rc - > tx_buf ) ; rc - > tx_buf = NULL ; rc - > tx_len = 0 ; }
if ( rc - > tx_buf ) { u_free ( rc - > tx_buf ) ; rc - > tx_buf = NULL ; rc - > tx_len = 0 ; }
if ( rc - > retry_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > retry_timer ) ; rc - > retry_timer = NULL ; }
if ( rc - > close_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > close_timer ) ; rc - > close_timer = NULL ; }
if ( rc - > close_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > close_timer ) ; rc - > close_timer = NULL ; }
if ( rc - > diag_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > diag_timer ) ; rc - > diag_timer = NULL ; }
if ( rc - > diag_timer ) { uasync_cancel_timeout ( rc - > ua , rc - > diag_timer ) ; rc - > diag_timer = NULL ; }
if ( rc - > tc ) { tcp_conn_destroy ( rc - > tc ) ; rc - > tc = NULL ; }
if ( rc - > tc ) { tcp_conn_destroy ( rc - > tc ) ; rc - > tc = NULL ; }