@ -17,7 +17,6 @@
# include <sys/time.h>
# include <sys/time.h>
# include <math.h> // For bandwidth calcs
# include <math.h> // For bandwidth calcs
# include <limits.h> // For UINT16_MAX
# include <limits.h> // For UINT16_MAX
# include <time.h> // For strftime in metrics snapshot
# include "../lib/mem.h"
# include "../lib/mem.h"
# include "../lib/strbuf.h"
# include "../lib/strbuf.h"
# include "../lib/memory_pool.h"
# include "../lib/memory_pool.h"
@ -405,7 +404,6 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
// Cancel active timers
// Cancel active timers
if ( etcp - > retrans_timer ) { uasync_cancel_timeout ( etcp - > instance - > ua , etcp - > retrans_timer ) ; etcp - > retrans_timer = NULL ; }
if ( etcp - > retrans_timer ) { uasync_cancel_timeout ( etcp - > instance - > ua , etcp - > retrans_timer ) ; etcp - > retrans_timer = NULL ; }
if ( etcp - > ack_resp_timer ) { uasync_cancel_timeout ( etcp - > instance - > ua , etcp - > ack_resp_timer ) ; etcp - > ack_resp_timer = NULL ; }
if ( etcp - > ack_resp_timer ) { uasync_cancel_timeout ( etcp - > instance - > ua , etcp - > ack_resp_timer ) ; etcp - > ack_resp_timer = NULL ; }
etcp_metrics_stop_timer ( etcp ) ;
routing_del_conn ( etcp ) ;
routing_del_conn ( etcp ) ;
etcp_router_transit_queues_destroy ( etcp ) ;
etcp_router_transit_queues_destroy ( etcp ) ;
@ -706,8 +704,6 @@ void etcp_conn_queue_set_ready(struct ETCP_CONN* conn) {
DEBUG_INFO ( DEBUG_CATEGORY_ETCP , " [%s] moved to ready queue, state=%d peer_node_id=0x%llx " ,
DEBUG_INFO ( DEBUG_CATEGORY_ETCP , " [%s] moved to ready queue, state=%d peer_node_id=0x%llx " ,
conn - > log_name , conn - > state , ( unsigned long long ) conn - > peer_node_id ) ;
conn - > log_name , conn - > state , ( unsigned long long ) conn - > peer_node_id ) ;
etcp_metrics_start_timer ( conn ) ;
etcp_cbk_fire ( conn , ETCP_CBK_EVENT_INIT ) ;
etcp_cbk_fire ( conn , ETCP_CBK_EVENT_INIT ) ;
if ( conn - > reinit_pending ) { conn - > reinit_pending = 0 ; etcp_cbk_fire ( conn , ETCP_CBK_EVENT_REINIT ) ; }
if ( conn - > reinit_pending ) { conn - > reinit_pending = 0 ; etcp_cbk_fire ( conn , ETCP_CBK_EVENT_REINIT ) ; }
if ( conn - > links_up ) etcp_on_up ( conn ) ;
if ( conn - > links_up ) etcp_on_up ( conn ) ;
@ -973,7 +969,6 @@ static void ack_timeout_check(struct ETCP_CONN* etcp) {
// Update stats
// Update stats
etcp - > retransmissions_count + + ;
etcp - > retransmissions_count + + ;
etcp_metrics_add_loss ( etcp , 1 ) ;
}
}
else { // не надо до конца сканировать - они уже сортированы по таймстемпу т.к. очередь fifo, а timestamp = время добавления в очередь = время отправки
else { // не надо до конца сканировать - они уже сортированы по таймстемпу т.к. очередь fifo, а timestamp = время добавления в очередь = время отправки
// shedule timer
// shedule timer
@ -1176,7 +1171,6 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
inf_pkt - > state = INFLIGHT_STATE_WAIT_ACK ;
inf_pkt - > state = INFLIGHT_STATE_WAIT_ACK ;
queue_data_put_with_index ( etcp - > input_wait_ack , & inf_pkt - > ll ) ; // move dgram to wait_ack queue
queue_data_put_with_index ( etcp - > input_wait_ack , & inf_pkt - > ll ) ; // move dgram to wait_ack queue
etcp_metrics_add_sent ( etcp , inf_pkt - > ll . len ) ;
}
}
size_t ack_q_size = queue_entry_count ( etcp - > ack_q ) ;
size_t ack_q_size = queue_entry_count ( etcp - > ack_q ) ;
@ -1582,7 +1576,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
pkt - > link - > rtt_last = new_rtt ;
pkt - > link - > rtt_last = new_rtt ;
DEBUG_INFO ( DEBUG_CATEGORY_DEBUG , " [%s] KA-RTT: rtt=%u cur=%u ret=%u dlen=%u " ,
DEBUG_INFO ( DEBUG_CATEGORY_DEBUG , " [%s] KA-RTT: rtt=%u cur=%u ret=%u dlen=%u " ,
etcp - > log_name , new_rtt , cur_ts , ret_ts , pkt - > data_len ) ;
etcp - > log_name , new_rtt , cur_ts , ret_ts , pkt - > data_len ) ;
etcp_metrics_add_rtt ( etcp , new_rtt ) ;
int recv_dt_tx1 = data [ 3 ] | ( data [ 4 ] < < 8 ) ; // localtime удаленной стороны момента принятия пакета - timestamp этого пакета (на стороне отправителя, т.е. у нас)
int recv_dt_tx1 = data [ 3 ] | ( data [ 4 ] < < 8 ) ; // localtime удаленной стороны момента принятия пакета - timestamp этого пакета (на стороне отправителя, т.е. у нас)
@ -1710,7 +1703,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
memcpy ( payload_data , data + 5 , pkt_len ) ;
memcpy ( payload_data , data + 5 , pkt_len ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] RX seq=%u need=%u asm_len=%d " , etcp - > log_name , seq , etcp - > last_delivered_id + 1 , etcp - > recv_q - > count ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] RX seq=%u need=%u asm_len=%d " , etcp - > log_name , seq , etcp - > last_delivered_id + 1 , etcp - > recv_q - > count ) ;
queue_data_put_with_index ( etcp - > recv_q , ( struct ll_entry * ) rx_pkt ) ;
queue_data_put_with_index ( etcp - > recv_q , ( struct ll_entry * ) rx_pkt ) ;
etcp_metrics_add_rcvd ( etcp , pkt_len ) ;
if ( ( int32_t ) ( seq - etcp - > last_delivered_id ) = = 1 ) etcp_output_try_assembly ( etcp ) ; // пробуем собрать выходную очередь из фрагментов
if ( ( int32_t ) ( seq - etcp - > last_delivered_id ) = = 1 ) etcp_output_try_assembly ( etcp ) ; // пробуем собрать выходную очередь из фрагментов
} else {
} else {
etcp - > rx_dup_count + + ;
etcp - > rx_dup_count + + ;
@ -1737,7 +1729,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
}
}
if ( link - > burst_recv_valid ) {
if ( link - > burst_recv_valid ) {
if ( seq ! = link - > burst_recv_next_seq ) {
if ( seq ! = link - > burst_recv_next_seq ) {
if ( seq > link - > burst_recv_next_seq ) etcp_metrics_add_loss ( etcp , seq - link - > burst_recv_next_seq ) ;
link - > burst_recv_valid = 0 ;
link - > burst_recv_valid = 0 ;
}
}
if ( seq < 16 ) {
if ( seq < 16 ) {
@ -1806,44 +1797,6 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
break ;
break ;
}
}
case ETCP_SECTION_METRICS : {
if ( len < 2 + SC_SIGN_SIZE ) { len = 0 ; break ; }
uint16_t csv_len = data [ 1 ] | ( data [ 2 ] < < 8 ) ;
if ( ( uint32_t ) ( 2 + SC_SIGN_SIZE + csv_len ) > len ) { len = 0 ; break ; }
uint8_t * sig = data + 2 ;
uint8_t * csv = data + 2 + SC_SIGN_SIZE ;
uint8_t * peer_pubkey = pkt - > link - > remote_ed25519_pubkey ;
struct sc_stream_sign_state verify_state ;
int sig_ok = ( sc_stream_sign_verify_init ( & verify_state , peer_pubkey ) = = SC_OK
& & sc_stream_sign_update ( & verify_state , csv , csv_len ) = = SC_OK
& & sc_stream_sign_verify ( & verify_state , sig , SC_SIGN_SIZE ) = = SC_OK ) ;
if ( ! sig_ok & & verify_state . initialized ) /* update failed, then verify didn't run */
sc_stream_sign_cleanup ( & verify_state ) ;
if ( sig_ok ) {
const char * stats_dir = etcp - > instance - > stats_dir ;
if ( stats_dir [ 0 ] & & etcp - > peer_node_id ! = 0 ) {
char path [ 768 ] ;
snprintf ( path , sizeof ( path ) , " %s/from_%016llx.dat " , stats_dir , ( unsigned long long ) etcp - > peer_node_id ) ;
FILE * f = fopen ( path , " a " ) ;
if ( f ) {
fprintf ( f , " %.*s, " , csv_len , csv ) ;
for ( int i = 0 ; i < SC_SIGN_SIZE ; i + + ) fprintf ( f , " %02x " , sig [ i ] ) ;
fprintf ( f , " \n " ) ;
fclose ( f ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] Received metrics from peer (%u bytes) " , etcp - > log_name , csv_len ) ;
if ( etcp - > instance - > api_bindings . on_metrics_rcvd )
etcp - > instance - > api_bindings . on_metrics_rcvd ( etcp - > instance - > api_bindings . metrics_user_ptr , etcp , csv , csv_len , sig ) ;
} else {
DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " [%s] Cannot open from_metrics file: %s " , etcp - > log_name , path ) ;
}
}
} else {
DEBUG_ERROR ( DEBUG_CATEGORY_CRYPTO , " [%s] Metrics signature verification FAILED " , etcp - > log_name ) ;
}
data + = 2 + SC_SIGN_SIZE + csv_len ; len - = 2 + SC_SIGN_SIZE + csv_len ;
break ;
}
default :
default :
DEBUG_WARN ( DEBUG_CATEGORY_ETCP , " unknown section type=0x%02x " , type ) ;
DEBUG_WARN ( DEBUG_CATEGORY_ETCP , " unknown section type=0x%02x " , type ) ;
len = 0 ;
len = 0 ;
@ -1883,178 +1836,4 @@ void etcp_update_mtu(struct ETCP_CONN* etcp) {
}
}
}
}
// ====================================================================== Метрики: гистограммы RTT/потерь + тотальные счётчики
static const uint16_t etcp_metrics_rtt_bounds [ ] = { 10 , 20 , 30 , 50 , 80 , 130 , 210 , 340 , 550 , 900 , 1500 } ;
# define ETCP_METRICS_RTT_BOUNDS_COUNT (sizeof(etcp_metrics_rtt_bounds) / sizeof(etcp_metrics_rtt_bounds[0]))
void etcp_metrics_init ( struct etcp_metrics * m ) {
if ( ! m ) return ;
memset ( m , 0 , sizeof ( * m ) ) ;
}
void etcp_metrics_add_rtt ( struct ETCP_CONN * etcp , uint16_t rtt_tb ) {
if ( ! etcp | | etcp - > state = = 2 ) return ;
struct etcp_metrics * m = & etcp - > metrics ;
int i ;
for ( i = 0 ; i < ( int ) ETCP_METRICS_RTT_BOUNDS_COUNT ; i + + ) {
if ( rtt_tb < etcp_metrics_rtt_bounds [ i ] ) break ;
}
if ( i > = ETCP_METRICS_RTT_BUCKETS ) i = ETCP_METRICS_RTT_BUCKETS - 1 ;
m - > rtt_hist [ i ] + + ;
m - > work_samples + + ;
}
void etcp_metrics_add_sent ( struct ETCP_CONN * etcp , uint32_t len ) {
if ( ! etcp | | etcp - > state = = 2 ) return ;
struct etcp_metrics * m = & etcp - > metrics ;
m - > work_sent + + ;
m - > work_bytes_sent + = len ;
m - > total_sent + + ;
m - > total_bytes_sent + = len ;
}
void etcp_metrics_add_rcvd ( struct ETCP_CONN * etcp , uint32_t len ) {
if ( ! etcp | | etcp - > state = = 2 ) return ;
struct etcp_metrics * m = & etcp - > metrics ;
m - > work_rcvd + + ;
m - > work_bytes_rcvd + = len ;
m - > total_rcvd + + ;
m - > total_bytes_rcvd + = len ;
}
void etcp_metrics_add_loss ( struct ETCP_CONN * etcp , uint32_t count ) {
if ( ! etcp | | etcp - > state = = 2 ) return ;
struct etcp_metrics * m = & etcp - > metrics ;
m - > work_lost + = count ;
m - > total_lost + = count ;
}
static void etcp_metrics_do_snapshot ( struct ETCP_CONN * etcp ) {
if ( ! etcp | | ! etcp - > instance ) return ;
struct etcp_metrics * m = & etcp - > metrics ;
// Копируем working → final
memcpy ( m - > rtt_hist_final , m - > rtt_hist , sizeof ( m - > rtt_hist_final ) ) ;
memcpy ( m - > loss_hist_final , m - > loss_hist , sizeof ( m - > loss_hist_final ) ) ;
m - > samples_final = m - > work_samples ;
m - > sent_final = m - > work_sent ;
m - > rcvd_final = m - > work_rcvd ;
m - > bytes_sent_final = m - > work_bytes_sent ;
m - > bytes_rcvd_final = m - > work_bytes_rcvd ;
m - > lost_final = m - > work_lost ;
// Рассчитываем loss-гистограмму
uint32_t total = m - > work_sent + m - > work_lost ;
if ( total > 0 ) {
uint32_t loss_pct = ( m - > work_lost * 100 + total / 2 ) / total ;
int bucket = ( int ) loss_pct ;
if ( bucket > = ETCP_METRICS_LOSS_BUCKETS ) bucket = ETCP_METRICS_LOSS_BUCKETS - 1 ;
m - > loss_hist [ bucket ] + + ;
}
// Обнуляем рабочую копию
memset ( m - > rtt_hist , 0 , sizeof ( m - > rtt_hist ) ) ;
memset ( m - > loss_hist , 0 , sizeof ( m - > loss_hist ) ) ;
m - > work_samples = 0 ;
m - > work_sent = 0 ; m - > work_rcvd = 0 ;
m - > work_bytes_sent = 0 ; m - > work_bytes_rcvd = 0 ;
m - > work_lost = 0 ;
m - > last_snapshot_tb = get_time_tb ( ) ;
// Строим CSV строку
char csv_buf [ 4096 ] ;
int pos = 0 ;
time_t now = time ( NULL ) ;
struct tm * tm = localtime ( & now ) ;
char date_str [ 16 ] , time_str [ 16 ] ;
if ( tm ) { strftime ( date_str , sizeof ( date_str ) , " %Y-%m-%d " , tm ) ; strftime ( time_str , sizeof ( time_str ) , " %H:%M:%S " , tm ) ; }
else { strcpy ( date_str , " ---- " ) ; strcpy ( time_str , " ---- " ) ; }
pos + = snprintf ( csv_buf + pos , sizeof ( csv_buf ) - pos , " %s,%s " , date_str , time_str ) ;
for ( int i = 0 ; i < ETCP_METRICS_RTT_BUCKETS ; i + + ) pos + = snprintf ( csv_buf + pos , sizeof ( csv_buf ) - pos , " ,%u " , m - > rtt_hist_final [ i ] ) ;
for ( int i = 0 ; i < ETCP_METRICS_LOSS_BUCKETS ; i + + ) pos + = snprintf ( csv_buf + pos , sizeof ( csv_buf ) - pos , " ,%u " , m - > loss_hist_final [ i ] ) ;
pos + = snprintf ( csv_buf + pos , sizeof ( csv_buf ) - pos , " ,%u,%u,%u,%u,%u,%u " , m - > samples_final , m - > sent_final , m - > rcvd_final , m - > bytes_sent_final , m - > bytes_rcvd_final , m - > lost_final ) ;
pos + = snprintf ( csv_buf + pos , sizeof ( csv_buf ) - pos , " ,%llu,%llu,%llu,%llu,%llu " , ( unsigned long long ) m - > total_sent , ( unsigned long long ) m - > total_rcvd , ( unsigned long long ) m - > total_bytes_sent , ( unsigned long long ) m - > total_bytes_rcvd , ( unsigned long long ) m - > total_lost ) ;
uint32_t csv_len = ( uint32_t ) pos ;
// Запись локального файла
const char * stats_dir = etcp - > instance - > stats_dir ;
if ( stats_dir [ 0 ] & & etcp - > peer_node_id ! = 0 ) {
char path [ 768 ] ;
snprintf ( path , sizeof ( path ) , " %s/%016llx.dat " , stats_dir , ( unsigned long long ) etcp - > peer_node_id ) ;
FILE * f = fopen ( path , " a " ) ;
if ( f ) { fprintf ( f , " %s \n " , csv_buf ) ; fclose ( f ) ; }
else DEBUG_ERROR ( DEBUG_CATEGORY_ETCP , " [%s] Cannot open metrics file: %s " , etcp - > log_name , path ) ;
}
DEBUG_INFO ( DEBUG_CATEGORY_ETCP , " [%s] Metrics snapshot: samples=%u sent=%u rcvd=%u lost=%u " , etcp - > log_name , m - > samples_final , m - > sent_final , m - > rcvd_final , m - > lost_final ) ;
// Подпись Ed25519 и отправка пиру
if ( etcp - > crypto_ctx . initialized & & etcp - > initialized ) {
struct sc_stream_sign_state sign_state ;
if ( sc_stream_sign_init ( & etcp - > crypto_ctx , & sign_state ) = = SC_OK ) {
uint8_t sig [ SC_SIGN_SIZE ] ;
size_t sig_len = sizeof ( sig ) ;
if ( sc_stream_sign_update ( & sign_state , ( uint8_t * ) csv_buf , csv_len ) = = SC_OK
& & sc_stream_sign_final ( & sign_state , sig , & sig_len ) = = SC_OK ) {
// Построить и отправить дграмму с секцией METRICS
struct ETCP_LINK * link = etcp_loadbalancer_select_link ( etcp ) ;
if ( link ) {
struct ETCP_DGRAM * dgram = memory_pool_alloc ( etcp - > instance - > pkt_pool ) ;
if ( dgram ) {
dgram - > link = link ;
dgram - > noencrypt_len = 0 ;
dgram - > timestamp = get_current_timestamp ( ) ;
int ptr = 0 ;
// Минимальная ACK-секция
dgram - > data [ ptr + + ] = 1 ; dgram - > data [ ptr + + ] = 0 ;
dgram - > data [ ptr + + ] = etcp - > last_delivered_id ; dgram - > data [ ptr + + ] = etcp - > last_delivered_id > > 8 ;
dgram - > data [ ptr + + ] = etcp - > last_delivered_id > > 16 ; dgram - > data [ ptr + + ] = etcp - > last_delivered_id > > 24 ;
dgram - > data [ ptr + + ] = etcp - > rx_dup_count ; dgram - > data [ ptr + + ] = etcp - > rx_dup_count > > 8 ;
// METRICS секция
dgram - > data [ ptr + + ] = ETCP_SECTION_METRICS ;
uint16_t csv16 = ( uint16_t ) csv_len ;
dgram - > data [ ptr + + ] = csv16 & 0xFF ; dgram - > data [ ptr + + ] = ( csv16 > > 8 ) & 0xFF ;
memcpy ( dgram - > data + ptr , sig , SC_SIGN_SIZE ) ; ptr + = SC_SIGN_SIZE ;
memcpy ( dgram - > data + ptr , csv_buf , csv_len ) ; ptr + = csv_len ;
dgram - > data_len = ptr ;
etcp_encrypt_send ( dgram ) ;
memory_pool_free ( etcp - > instance - > pkt_pool , dgram ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] Metrics sent to peer (%u bytes) " , etcp - > log_name , csv_len ) ;
}
}
} else {
sc_stream_sign_cleanup ( & sign_state ) ;
}
}
}
}
static void metrics_snapshot_timer_cb ( void * arg ) {
struct ETCP_CONN * etcp = ( struct ETCP_CONN * ) arg ;
if ( ! etcp ) return ;
etcp - > metrics . timer = NULL ;
if ( etcp - > metrics . work_sent + etcp - > metrics . work_rcvd > = ETCP_METRICS_MIN_PACKETS )
etcp_metrics_do_snapshot ( etcp ) ;
// Перезапускаем таймер
etcp - > metrics . timer = uasync_set_timeout ( etcp - > instance - > ua , ETCP_METRICS_INTERVAL_TB , etcp , metrics_snapshot_timer_cb , " etcp_metrics " ) ;
}
void etcp_metrics_start_timer ( struct ETCP_CONN * etcp ) {
if ( ! etcp | | etcp - > state = = 2 | | etcp - > metrics . timer ) return ;
etcp_metrics_init ( & etcp - > metrics ) ;
etcp - > metrics . timer = uasync_set_timeout ( etcp - > instance - > ua , ETCP_METRICS_INTERVAL_TB , etcp , metrics_snapshot_timer_cb , " etcp_metrics " ) ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] Metrics timer started (interval=%ums) " , etcp - > log_name , ETCP_METRICS_INTERVAL_TB / 10 ) ;
}
void etcp_metrics_stop_timer ( struct ETCP_CONN * etcp ) {
if ( ! etcp ) return ;
if ( etcp - > metrics . timer & & etcp - > instance ) {
uasync_cancel_timeout ( etcp - > instance - > ua , etcp - > metrics . timer ) ;
etcp - > metrics . timer = NULL ;
DEBUG_DEBUG ( DEBUG_CATEGORY_ETCP , " [%s] Metrics timer stopped " , etcp - > log_name ) ;
}
}