@ -55,26 +55,28 @@ struct NODE_CONN_DIRECT {
struct NODE_CONN_DIRECT * next ;
} ;
static struct ncd_entry * g_ncd_registry ;
static uint8_t g_ncd_control_bound ; /* 1 = etcp_bind(ETCP_RT_ID_NCD_CONTROL) уже сделан */
/* ─── forward declarations ─── */
static void ncd_event_dispatch ( struct ncd_entry * entry , enum ncd_event event ) ;
static void ncd_connect_timeout_cb ( void * arg ) ;
static void ncd_deliver_up_cb ( void * arg ) ;
static void ncd_init_cb ( struct ETCP_CONN * conn , int event , void * arg ) ;
static void ncd_up_cb ( struct ETCP_CONN * conn , int event , void * arg ) ;
static void ncd_down_cb ( struct ETCP_CONN * conn , int event , void * arg ) ;
static void ncd_deferred_close ( void * arg ) ;
static void ncd_deferred_close_conn ( void * arg ) ;
/* ═══════════ реестр ═══════════ */
static struct ncd_entry * ncd_registry_find ( uint64_t node_id ) {
struct ncd_entry * e = g_ncd_registry ;
static struct ncd_entry * ncd_registry_find ( struct UTUN_INSTANCE * inst , uint64_t node_id ) {
struct ncd_entry * e = ( struct ncd_entry * ) inst - > ncd_registry ;
while ( e ) { if ( e - > node_id = = node_id ) return e ; e = e - > next ; }
return NULL ;
}
static void ncd_registry_add ( struct ncd_entry * entry ) {
entry - > next = g_ncd_registry ; g_ ncd_registry = entry ;
static void ncd_registry_add ( struct UTUN_INSTANCE * inst , struct ncd_entry * entry ) {
entry - > next = ( struct ncd_entry * ) inst - > ncd_registry ; inst - > ncd_registry = entry ;
}
static void ncd_registry_remove ( struct ncd_entry * entry ) {
struct ncd_entry * * pp = & g_ ncd_registry;
static void ncd_registry_remove ( struct UTUN_INSTANCE * inst , struct ncd_entry * entry ) {
struct ncd_entry * * pp = ( struct ncd_entry * * ) & inst - > ncd_registry ;
while ( * pp ) { if ( * pp = = entry ) { * pp = entry - > next ; return ; } pp = & ( * pp ) - > next ; }
}
@ -98,7 +100,8 @@ static struct TOPO_NODE* ncd_lookup_node(struct UTUN_INSTANCE* inst, uint64_t no
/* ═══════════ создание линков (round‑robin, все сокеты кроме PRIVATE) ═══════════ */
static int ncd_create_links ( struct ncd_entry * entry , struct TOPO_NODE * ni ) {
static int ncd_create_links ( struct ncd_entry * entry , struct TOPO_NODE * ni ,
struct ETCP_SOCKET * specific_sock ) {
struct ETCP_CONN * conn = entry - > conn ;
struct UTUN_INSTANCE * inst = conn - > instance ;
int link_count = 0 ;
@ -106,9 +109,14 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni) {
/* ---- IPv4 ---- */
{
struct ETCP_SOCKET * socks [ 64 ] ; int sock_count = 0 ;
if ( specific_sock ) {
if ( specific_sock - > local_addr . ss_family = = AF_INET )
socks [ sock_count + + ] = specific_sock ;
} else {
for ( struct ETCP_SOCKET * s = inst - > etcp_sockets ; s ; s = s - > next )
if ( s - > local_addr . ss_family = = AF_INET & & s - > type ! = CFG_SERVER_TYPE_PRIVATE )
socks [ sock_count + + ] = s ;
}
if ( sock_count > 0 ) {
int rr = 0 ;
@ -127,9 +135,14 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni) {
/* ---- IPv6 ---- */
{
struct ETCP_SOCKET * socks [ 64 ] ; int sock_count = 0 ;
if ( specific_sock ) {
if ( specific_sock - > local_addr . ss_family = = AF_INET6 )
socks [ sock_count + + ] = specific_sock ;
} else {
for ( struct ETCP_SOCKET * s = inst - > etcp_sockets ; s ; s = s - > next )
if ( s - > local_addr . ss_family = = AF_INET6 & & s - > type ! = CFG_SERVER_TYPE_PRIVATE )
socks [ sock_count + + ] = s ;
}
if ( sock_count > 0 ) {
int rr = 0 ;
@ -186,19 +199,32 @@ static void ncd_up_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)even
if ( ! entry | | conn - > fin_wait ) return ;
ncd_event_dispatch ( entry , NCD_EVENT_UP ) ;
}
static void ncd_deferred_close ( void * arg ) {
struct ncd_entry * entry = ( struct ncd_entry * ) arg ;
struct ETCP_CONN * conn = entry - > conn ;
if ( conn ) {
etcp_conn_remove_cbk ( conn , ncd_init_cb , entry ) ;
etcp_conn_remove_cbk ( conn , ncd_up_cb , entry ) ;
etcp_conn_remove_cbk ( conn , ncd_down_cb , entry ) ;
ncd_registry_remove ( conn - > instance , entry ) ;
etcp_connection_close ( conn ) ;
}
u_free ( entry ) ;
}
static void ncd_deferred_close_conn ( void * arg ) {
struct ETCP_CONN * conn = ( struct ETCP_CONN * ) arg ;
if ( conn ) etcp_connection_close ( conn ) ;
}
static void ncd_down_cb ( struct ETCP_CONN * conn , int event , void * arg ) { ( void ) event ;
struct ncd_entry * entry = ( struct ncd_entry * ) arg ;
if ( ! entry ) return ;
DEBUG_DEBUG ( DEBUG_CATEGORY_DEBUG , " [ncd-debug] ncd_down_cb conn=%p entry=%p fin_wait=%d handle_count=%d node=0x%016llx " ,
( void * ) conn , entry , conn - > fin_wait , entry - > handle_count , ( unsigned long long ) entry - > node_id ) ;
if ( conn - > fin_wait & & entry - > handle_count < = 0 ) {
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] DOWN during fin_wait, cleaning up node=0x%016llx " , ( unsigned long long ) entry - > node_id ) ;
conn - > fin_wait = 0 ; conn - > fin_wait_clear_cb = NULL ; conn - > fin_wait_clear_arg = NULL ;
if ( entry - > fin_wait_timer ) { uasync_cancel_timeout ( entry - > ua , entry - > fin_wait_timer ) ; entry - > fin_wait_timer = NULL ; }
etcp_conn_remove_cbk ( conn , ncd_init_cb , entry ) ;
etcp_conn_remove_cbk ( conn , ncd_up_cb , entry ) ;
etcp_conn_remove_cbk ( conn , ncd_down_cb , entry ) ;
etcp_connection_close ( conn ) ;
ncd_registry_remove ( entry ) ;
u_free ( entry ) ;
uasync_call_soon ( entry - > ua , entry , ncd_deferred_close ) ;
return ;
}
ncd_event_dispatch ( entry , NCD_EVENT_DOWN ) ;
@ -241,7 +267,7 @@ static void ncd_fin_wait_timeout_cb(void* arg) {
etcp_conn_remove_cbk ( entry - > conn , ncd_up_cb , entry ) ;
etcp_conn_remove_cbk ( entry - > conn , ncd_down_cb , entry ) ;
etcp_connection_close ( entry - > conn ) ;
ncd_registry_remove ( entry ) ;
ncd_registry_remove ( entry - > conn - > instance , entry ) ;
u_free ( entry ) ;
}
@ -254,14 +280,14 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e)
}
struct ncd_control_msg * msg = ( struct ncd_control_msg * ) e - > dgram ;
uint64_t sender_id = msg - > node_id ;
struct ncd_entry * entry = ncd_registry_find ( sender_id ) ;
struct ncd_entry * entry = ncd_registry_find ( conn - > instance , sender_id ) ;
switch ( msg - > subcmd ) {
case NCD_SUBCMD_CLOSE :
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] recv CLOSE from 0x%016llx " , ( unsigned long long ) sender_id ) ;
if ( ! entry ) {
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] no entry for CLOSE sender 0x%016llx, closing conn " , ( unsigned long long ) sender_id ) ;
if ( conn ) etcp_connection_close ( conn ) ;
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] no entry for CLOSE sender 0x%016llx, deferred close conn " , ( unsigned long long ) sender_id ) ;
if ( conn ) uasync_call_soon ( conn - > instance - > ua , conn , ncd_deferred_close_ conn) ;
break ;
}
if ( entry - > handle_count > 0 ) {
@ -274,8 +300,8 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e)
etcp_send ( conn , qe ) ; }
}
} else {
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] closing conn for CLOSE from 0x%016llx " , ( unsigned long long ) sender_id ) ;
if ( conn ) etcp_connection_close ( conn ) ;
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] deferred close conn for CLOSE from 0x%016llx " , ( unsigned long long ) sender_id ) ;
if ( conn ) uasync_call_soon ( conn - > instance - > ua , conn , ncd_deferred_close_ conn) ;
}
break ;
case NCD_SUBCMD_KEEP_ALIVE :
@ -283,8 +309,14 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e)
if ( entry & & entry - > conn - > fin_wait ) {
entry - > conn - > fin_wait = 0 ; entry - > conn - > fin_wait_clear_cb = NULL ; entry - > conn - > fin_wait_clear_arg = NULL ;
if ( entry - > fin_wait_timer ) { uasync_cancel_timeout ( entry - > ua , entry - > fin_wait_timer ) ; entry - > fin_wait_timer = NULL ; }
if ( entry - > handle_count < = 0 ) {
entry - > conn - > fin_wait = 1 ;
entry - > fin_wait_timer = uasync_set_timeout ( entry - > ua , NCD_FIN_WAIT_TIMEOUT_TB , entry , ncd_fin_wait_timeout_cb , " ncd_fin_wait " ) ;
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] KEEP_ALIVE but no handles, restart fin_wait " ) ;
} else {
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] KEEP_ALIVE: fin_wait cleared, conn stays alive " ) ;
}
}
break ;
}
queue_dgram_free ( e ) ; queue_entry_free ( e ) ;
@ -293,9 +325,9 @@ static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e)
/* ═══════════ инициализация глобального обработчика ═══════════ */
static void ncd_init_control_binding ( struct UTUN_INSTANCE * inst ) {
if ( g_ ncd_control_bound) return ;
if ( inst - > ncd_control_bound ) return ;
etcp_bind ( inst , ETCP_RT_ID_NCD_CONTROL , ncd_recv_control_handler ) ;
g_ ncd_control_bound = 1 ;
inst - > ncd_control_bound = 1 ;
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] bound ETCP_RT_ID_NCD_CONTROL handler " ) ;
}
@ -303,13 +335,14 @@ static void ncd_init_control_binding(struct UTUN_INSTANCE* inst) {
int node_conn_direct_open ( struct UTUN_INSTANCE * inst , uint64_t node_id ,
ncd_callback cb , void * cb_arg ,
struct NODE_CONN_DIRECT * * out_handle ) {
struct NODE_CONN_DIRECT * * out_handle ,
struct ETCP_SOCKET * specific_sock ) {
if ( ! inst | | ! out_handle ) return NCD_ERR ;
* out_handle = NULL ;
ncd_init_control_binding ( inst ) ;
/* 1. Ищем в реестре — conn уже есть */
struct ncd_entry * entry = ncd_registry_find ( node_id ) ;
struct ncd_entry * entry = ncd_registry_find ( inst , node_id ) ;
if ( entry ) {
/* снять fin_wait если был */
if ( entry - > conn & & entry - > conn - > fin_wait ) {
@ -344,11 +377,11 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
if ( ! entry ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] alloc entry failed node=0x%016llx " , ( unsigned long long ) node_id ) ; return NCD_ERR ; }
entry - > node_id = node_id ; entry - > conn = conn ; entry - > ua = inst - > ua ;
entry - > up = ( uint8_t ) ( conn - > links_up ? 1 : 0 ) ;
ncd_registry_add ( entry ) ;
ncd_registry_add ( inst , entry ) ;
struct NODE_CONN_DIRECT * h = u_calloc ( 1 , sizeof ( * h ) ) ;
if ( ! h ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] alloc handle failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
ncd_registry_remove ( entry ) ; u_free ( entry ) ; return NCD_ERR ; }
ncd_registry_remove ( inst , entry ) ; u_free ( entry ) ; return NCD_ERR ; }
h - > entry = entry ; h - > cb = cb ; h - > cb_arg = cb_arg ;
h - > next = entry - > handles ; entry - > handles = h ;
entry - > handle_count = 1 ;
@ -375,7 +408,11 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
return NCD_ERR ;
}
conn = etcp_connection_create ( inst , NULL ) ;
{ char conn_name [ MAX_CONN_NAME_LEN ] ;
if ( ni - > node_name & & ni - > node_name [ 0 ] ) strncpy ( conn_name , ni - > node_name , MAX_CONN_NAME_LEN - 1 ) ;
else snprintf ( conn_name , sizeof ( conn_name ) , " n_%016llx " , ( unsigned long long ) node_id ) ;
conn = etcp_connection_create ( inst , conn_name ) ; }
if ( ! conn ) {
DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] etcp_connection_create failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
topo_node_registry_unref ( inst - > topo_groups , ni - > node_id ) ;
@ -411,11 +448,11 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
etcp_connection_close ( conn ) ; topo_node_registry_unref ( inst - > topo_groups , ni - > node_id ) ; return NCD_ERR ;
}
entry - > node_id = node_id ; entry - > conn = conn ; entry - > ua = inst - > ua ; entry - > up = 0 ;
ncd_registry_add ( entry ) ;
ncd_registry_add ( inst , entry ) ;
struct NODE_CONN_DIRECT * h = u_calloc ( 1 , sizeof ( * h ) ) ;
if ( ! h ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] alloc handle failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
ncd_registry_remove ( entry ) ; u_free ( entry ) ; etcp_connection_close ( conn ) ; topo_node_registry_unref ( inst - > topo_groups , ni - > node_id ) ; return NCD_ERR ; }
ncd_registry_remove ( inst , entry ) ; u_free ( entry ) ; etcp_connection_close ( conn ) ; topo_node_registry_unref ( inst - > topo_groups , ni - > node_id ) ; return NCD_ERR ; }
h - > entry = entry ; h - > cb = cb ; h - > cb_arg = cb_arg ;
h - > next = entry - > handles ; entry - > handles = h ;
entry - > handle_count = 1 ;
@ -425,7 +462,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
etcp_conn_add_cbk ( conn , ncd_up_cb , entry , ETCP_CBK_EVENT_UP ) ;
etcp_conn_add_cbk ( conn , ncd_down_cb , entry , ETCP_CBK_EVENT_DOWN ) ;
int link_count = ncd_create_links ( entry , ni ) ;
int link_count = ncd_create_links ( entry , ni , specific_sock ) ;
if ( link_count = = 0 )
DEBUG_WARN ( DEBUG_CATEGORY_NCD , " [ncd] no links created for node=0x%016llx " , ( unsigned long long ) node_id ) ;
@ -440,13 +477,14 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
int node_conn_direct_open_node ( struct UTUN_INSTANCE * inst , uint64_t node_id ,
ncd_callback cb , void * cb_arg ,
struct NODE_CONN_DIRECT * * out_handle ,
struct TOPO_NODE * ni ) {
struct TOPO_NODE * ni ,
struct ETCP_SOCKET * specific_sock ) {
if ( ! inst | | ! out_handle | | ! ni ) return NCD_ERR ;
* out_handle = NULL ;
ncd_init_control_binding ( inst ) ;
/* 1. Ищем в реестре — conn уже есть */
{ struct ncd_entry * entry = ncd_registry_find ( node_id ) ;
{ struct ncd_entry * entry = ncd_registry_find ( inst , node_id ) ;
if ( entry ) {
if ( entry - > conn & & entry - > conn - > fin_wait ) {
DEBUG_INFO ( DEBUG_CATEGORY_NCD , " [ncd] open_node REUSED clearing fin_wait node=0x%016llx " , ( unsigned long long ) node_id ) ;
@ -479,11 +517,11 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
if ( ! entry ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] open_node alloc entry failed node=0x%016llx " , ( unsigned long long ) node_id ) ; return NCD_ERR ; }
entry - > node_id = node_id ; entry - > conn = conn ; entry - > ua = inst - > ua ;
entry - > up = ( uint8_t ) ( conn - > links_up ? 1 : 0 ) ;
ncd_registry_add ( entry ) ;
ncd_registry_add ( inst , entry ) ;
struct NODE_CONN_DIRECT * h = u_calloc ( 1 , sizeof ( * h ) ) ;
if ( ! h ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] open_node alloc handle failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
ncd_registry_remove ( entry ) ; u_free ( entry ) ; return NCD_ERR ; }
ncd_registry_remove ( inst , entry ) ; u_free ( entry ) ; return NCD_ERR ; }
h - > entry = entry ; h - > cb = cb ; h - > cb_arg = cb_arg ;
h - > next = entry - > handles ; entry - > handles = h ;
entry - > handle_count = 1 ;
@ -504,7 +542,10 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
} }
/* 3. Новое подключение — используем переданный ni (временный, не владеем) */
{ struct ETCP_CONN * conn = etcp_connection_create ( inst , NULL ) ;
{ char conn_name [ MAX_CONN_NAME_LEN ] ;
if ( ni - > node_name & & ni - > node_name [ 0 ] ) strncpy ( conn_name , ni - > node_name , MAX_CONN_NAME_LEN - 1 ) ;
else snprintf ( conn_name , sizeof ( conn_name ) , " n_%016llx " , ( unsigned long long ) node_id ) ;
struct ETCP_CONN * conn = etcp_connection_create ( inst , conn_name ) ;
if ( ! conn ) {
DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] open_node etcp_connection_create failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
return NCD_ERR ;
@ -538,11 +579,11 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
etcp_connection_close ( conn ) ; return NCD_ERR ;
}
entry - > node_id = node_id ; entry - > conn = conn ; entry - > ua = inst - > ua ; entry - > up = 0 ;
ncd_registry_add ( entry ) ;
ncd_registry_add ( inst , entry ) ;
struct NODE_CONN_DIRECT * h = u_calloc ( 1 , sizeof ( * h ) ) ;
if ( ! h ) { DEBUG_ERROR ( DEBUG_CATEGORY_NCD , " [ncd] open_node alloc handle failed node=0x%016llx " , ( unsigned long long ) node_id ) ;
ncd_registry_remove ( entry ) ; u_free ( entry ) ; etcp_connection_close ( conn ) ; return NCD_ERR ; }
ncd_registry_remove ( inst , entry ) ; u_free ( entry ) ; etcp_connection_close ( conn ) ; return NCD_ERR ; }
h - > entry = entry ; h - > cb = cb ; h - > cb_arg = cb_arg ;
h - > next = entry - > handles ; entry - > handles = h ;
entry - > handle_count = 1 ;
@ -552,7 +593,7 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
etcp_conn_add_cbk ( conn , ncd_up_cb , entry , ETCP_CBK_EVENT_UP ) ;
etcp_conn_add_cbk ( conn , ncd_down_cb , entry , ETCP_CBK_EVENT_DOWN ) ;
int link_count = ncd_create_links ( entry , ni ) ;
int link_count = ncd_create_links ( entry , ni , specific_sock ) ;
if ( link_count = = 0 )
DEBUG_WARN ( DEBUG_CATEGORY_NCD , " [ncd] open_node no links created for node=0x%016llx " , ( unsigned long long ) node_id ) ;
@ -601,9 +642,11 @@ void node_conn_direct_close(struct NODE_CONN_DIRECT* h) {
etcp_conn_remove_cbk ( conn , ncd_up_cb , entry ) ;
etcp_conn_remove_cbk ( conn , ncd_down_cb , entry ) ;
etcp_connection_close ( conn ) ;
}
ncd_registry_remove ( entry ) ;
ncd_registry_remove ( conn - > instance , entry ) ;
u_free ( entry ) ;
} else {
DEBUG_WARN ( DEBUG_CATEGORY_NCD , " [ncd] close: no conn for entry, leaking entry=%p " , entry ) ;
}
}
}