@ -32,6 +32,7 @@ static void on_flushed_cb(struct tcp_conn* tc, void* arg);
static void read_queue_drain_cb ( struct ll_queue * q , void * arg ) ;
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 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 ) ;
uint32_t sid , const uint8_t * data , size_t len ) ;
static void send_close ( struct tcp_proxy_server_conn * rc ) ;
static void send_close ( struct tcp_proxy_server_conn * rc ) ;
@ -123,6 +124,35 @@ static void on_error_cb(struct tcp_conn* tc, int err, void* arg) {
tcp_proxy_server_conn_free ( rc ) ;
tcp_proxy_server_conn_free ( rc ) ;
}
}
static void diag_timer_cb ( void * arg ) {
struct tcp_proxy_server_conn * rc = ( struct tcp_proxy_server_conn * ) arg ;
struct tcp_conn * tc = rc - > tc ;
if ( ! tc | | tc - > sock = = SOCKET_INVALID ) return ;
size_t ea , er , da , dr ;
memory_pool_get_stats ( tc - > entry_pool , & ea , & er ) ;
memory_pool_get_stats ( tc - > data_pool , & da , & dr ) ;
int rcv_buf = 0 , snd_buf = 0 ;
socklen_t optlen = sizeof ( int ) ;
getsockopt ( tc - > sock , SOL_SOCKET , SO_RCVBUF , & rcv_buf , & optlen ) ;
getsockopt ( tc - > sock , SOL_SOCKET , SO_SNDBUF , & snd_buf , & optlen ) ;
DEBUG_INFO ( DEBUG_CATEGORY_SOCKET ,
" SOCK:DIAG fd=%d sid=%08x conn=%d err=%d fin=%d "
" rpaused=%d wmon=%d rq=%d(%zdb) wq=%d(%zdb) wbuf=%s "
" entry=%zd data=%zd tcp_rb=%d tcp_sb=%d " ,
( int ) tc - > sock , rc - > stream_id , tc - > connected , tc - > error , tc - > fin ,
tc - > read_paused , tc - > write_monitor ,
tc - > read_queue - > count , queue_total_bytes ( tc - > read_queue ) ,
tc - > write_queue - > count , queue_total_bytes ( tc - > write_queue ) ,
tc - > write_buf ? " y " : " n " ,
( ssize_t ) ( ea - er ) , ( ssize_t ) ( da - dr ) ,
rcv_buf , snd_buf ) ;
rc - > diag_timer = uasync_set_timeout ( rc - > ua , 10000 , rc , diag_timer_cb , " tps_diag " ) ;
}
// ====================================================================
// ====================================================================
// Дрейн read_queue → ETCP (автозабор + backpressure)
// Дрейн read_queue → ETCP (автозабор + backpressure)
// ====================================================================
// ====================================================================
@ -173,6 +203,7 @@ void tcp_proxy_server_conn_free(struct tcp_proxy_server_conn* rc) {
}
}
if ( rc - > tc ) { tcp_conn_destroy ( rc - > tc ) ; rc - > tc = NULL ; }
if ( rc - > tc ) { tcp_conn_destroy ( rc - > tc ) ; rc - > tc = 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 - > ctx & & rc - > ctx - > inst ) etcp_router_waiter_cancel ( rc - > ctx - > inst , rc - > peer_node_id , & rc - > pause_waiter ) ;
if ( rc - > ctx & & rc - > ctx - > inst ) etcp_router_waiter_cancel ( rc - > ctx - > inst , rc - > peer_node_id , & rc - > pause_waiter ) ;
u_free ( rc ) ;
u_free ( rc ) ;
}
}
@ -224,6 +255,7 @@ int tcp_proxy_server_handle_connect(struct UTUN_INSTANCE* inst, struct ll_entry*
}
}
rc - > next = ctx - > conns ; ctx - > conns = rc ;
rc - > next = ctx - > conns ; ctx - > conns = rc ;
rc - > diag_timer = uasync_set_timeout ( rc - > ua , 10000 , rc , diag_timer_cb , " tps_diag " ) ;
queue_dgram_free ( entry ) ; queue_entry_free ( entry ) ;
queue_dgram_free ( entry ) ; queue_entry_free ( entry ) ;
return 0 ;
return 0 ;
}
}