diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index fcf7e4b4..64713fc1 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -18,6 +18,7 @@ #include "stcp_link.h" #include "topo_node.h" #include "topo_group.h" +#include "node_conn_direct.h" #include "nat_detection.h" #include "../lib/memory_pool.h" #include "../lib/u_async.h" @@ -2200,137 +2201,73 @@ int init_connections(struct UTUN_INSTANCE* instance) { } } - // Initialize clients - create outgoing connections + // Initialize clients via node_conn_direct struct CFG_CLIENT* client = config->clients; - DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, total_conns=%d", - instance, config, config ? config->clients : NULL, instance ? queue_entry_count(instance->connections) : -1); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections: clients=%p total_conns=%d", config->clients, queue_entry_count(instance->connections)); while (client) { - // Check if client has required configuration - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Client %s - keepalive=%d, links=%p, peer_key_len=%zu", - client->name, client->keepalive, client->links, - strlen(client->peer_public_key_hex)); - - // Create ETCP connection for this client - struct ETCP_CONN* etcp_conn = etcp_connection_create(instance, client->name); - - if (!etcp_conn) { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to create ETCP connection for client %s", client->name); - client = client->next; - continue; - } - - // Generate session_id for this client connection - if (random_bytes((uint8_t*)&etcp_conn->session_id, sizeof(etcp_conn->session_id)) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "Failed to generate session_id for client %s", client->name); - etcp_connection_close(etcp_conn); - client = client->next; - continue; - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Client %s session_id=%08x", client->name, etcp_conn->session_id); - - // Initialize crypto context for this connection - if (sc_init_ctx(&etcp_conn->crypto_ctx, &instance->my_keys) != SC_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to initialize crypto context for client %s", client->name); - etcp_connection_close(etcp_conn); - client = client->next; - continue; - } - // If client has peer public key configured, set it - if (strlen(client->peer_public_key_hex) > 0) { - // For now, set peer node ID to indicate we have peer key - // The actual peer key will be exchanged during connection establishment - etcp_conn->peer_node_id = 1; // Simple indicator - etcp_update_log_name(etcp_conn); // Update log_name with peer_node_id - - DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "setting peer public key for client %s", client->name); - // Set peer public key (assuming hex format) - if (sc_set_peer_public_key(&etcp_conn->crypto_ctx, (const uint8_t*)client->peer_public_key_hex, 1) != SC_OK) { - DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to set peer public key for client %s", client->name); - } else { - DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "successfully set peer public key for client %s", client->name); - } - } else { + if (strlen(client->peer_public_key_hex) == 0) { DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "no peer public key configured for client %s", client->name); + client = client->next; continue; } - - etcp_set_routing_exchange_state(etcp_conn, 1); // инициируем обмен маршрутами - - // Create links for this client - struct CFG_CLIENT_LINK* client_link = client->links; - while (client_link) { - // Find the local server for this link - struct CFG_SERVER* local_server = client_link->local_srv; - if (!local_server) { - client_link = client_link->next; + uint8_t pubkey_bin[SC_PUBKEY_SIZE]; + if (sc_hex_to_binary(client->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "invalid peer pubkey hex for client %s", client->name); + client = client->next; continue; + } + uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "client %s node_id=0x%016llx", client->name, (unsigned long long)node_id); + + struct NODE_CONN_DIRECT* handle = NULL; + struct ETCP_CONN* conn = NULL; + for (struct CFG_CLIENT_LINK* cl = client->links; cl; cl = cl->next) { + if (cl->local_srv && cl->local_srv->transport) { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "client %s TCP transport not yet supported via NCD, skipping", client->name); continue; } - - // Find the socket for this server - struct ETCP_SOCKET* e_sock = NULL; - struct ETCP_SOCKET* sock = instance->etcp_sockets; - while (sock) { - if (sock->local_addr.ss_family == local_server->ip.ss_family) { - if (sock->local_addr.ss_family == AF_INET) { - struct sockaddr_in* sock_addr = (struct sockaddr_in*)&sock->local_addr; - struct sockaddr_in* srv_addr = (struct sockaddr_in*)&local_server->ip; - if (sock_addr->sin_addr.s_addr == srv_addr->sin_addr.s_addr && - sock_addr->sin_port == srv_addr->sin_port) { - e_sock = sock; - break; - } - } else if (sock->local_addr.ss_family == AF_INET6) { - struct sockaddr_in6* sock_addr6 = (struct sockaddr_in6*)&sock->local_addr; - struct sockaddr_in6* srv_addr6 = (struct sockaddr_in6*)&local_server->ip; - if (memcmp(&sock_addr6->sin6_addr, &srv_addr6->sin6_addr, 16) == 0 && - sock_addr6->sin6_port == srv_addr6->sin6_port) { - e_sock = sock; - break; - } - } - } - sock = sock->next; + struct ETCP_SOCKET* sock = NULL; + for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next) + if (cl->local_srv && strcmp(cl->local_srv->name, s->name) == 0) { sock = s; break; } + if (!sock) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s link", + cl->local_srv ? cl->local_srv->name : "?", client->name); + continue; } - - if (local_server->transport) { - if (strlen(client->peer_public_key_hex) == 0) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TCP client %s has no peer key", client->name); - client_link = client_link->next; continue; + if (!handle) { + struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr; + struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp)); + ni_tmp.node_id = node_id; ni_tmp.node_name = client->name; + memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE); + if (cl->remote_addr.ss_family == AF_INET) { + struct sockaddr_in* sin = (struct sockaddr_in*)&cl->remote_addr; + v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4); + v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP; + ni_tmp.v4_addrs = &v4_addr; + } else if (cl->remote_addr.ss_family == AF_INET6) { + struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&cl->remote_addr; + v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16); + v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP; + ni_tmp.v6_addrs = &v6_addr; } - // TCP transport — use stcp_link instead of etcp_link_new - uint16_t rport = ntohs(((struct sockaddr_in*)&client_link->remote_addr)->sin_port); - struct stcp_link_config scfg = { - .ua = instance->ua, .my_keys = &instance->my_keys, .inst = instance, - .peer_pubkey = (const uint8_t*)client->peer_public_key_hex, - .peer_pubkey_mode = 1, // hex - .remote_addr = &client_link->remote_addr, .remote_port = rport - }; - struct stcp_link *slink = stcp_link_connect(&scfg); - if (!slink) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create TCP link for client %s", client->name); - client_link = client_link->next; continue; + int r = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock); + if (r == NCD_ERR || !handle) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd open failed for client %s", client->name); + break; } - etcp_conn->transport_link = slink; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP link created for client %s", client->name); - client_link = client_link->next; continue; + conn = node_conn_direct_get_conn(handle); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "client %s ncd handle=%p conn=%p %s sock=%s", + client->name, handle, conn, r == NCD_NEW ? "NEW" : "REUSED", sock->name); + } else { + etcp_link_new(conn, sock, &cl->remote_addr, 0); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "client %s added link sock=%s", client->name, sock->name); } + } - if (!e_sock) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "No socket found for client %s link", client->name); - client_link = client_link->next; - continue; - } - - // Create link for this client connection - struct ETCP_LINK* link = etcp_link_new(etcp_conn, e_sock, &client_link->remote_addr, 0); // 0 = client initiates - if (!link) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create link for client %s", client->name); - client_link = client_link->next; - continue; - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Created link %p for client %s, socket=%p", - link, client->name, e_sock); - - client_link = client_link->next; + if (handle) { + struct CONFIG_CONN_HANDLE* ch = u_calloc(1, sizeof(*ch)); + if (ch) { ch->node_id = node_id; strncpy(ch->name, client->name, MAX_CONN_NAME_LEN - 1); + ch->handle = handle; ch->next = instance->config_conn_handles; + instance->config_conn_handles = ch; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "client %s saved handle to config_conn_handles", client->name); } } client = client->next; diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index 644009f3..2d0e4558 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/src/transport_layer/node_conn_direct.c @@ -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; - 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 (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; - 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 (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,7 +309,13 @@ 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; } - DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] KEEP_ALIVE: fin_wait cleared, conn stays alive"); + 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; } @@ -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(conn->instance, entry); + u_free(entry); + } else { + DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] close: no conn for entry, leaking entry=%p", entry); } - ncd_registry_remove(entry); - u_free(entry); } } diff --git a/src/utun_instance.c b/src/utun_instance.c index 3e236c97..df6194e2 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -18,6 +18,7 @@ #include "chat/chat_sync.h" #include "stcp_server.h" #include "control_server.h" +#include "transport_layer/node_conn_direct.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" @@ -418,6 +419,17 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { instance->control_srv = NULL; } + // Close config-based connection handles (sends CLOSE before sockets die) + struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles; + while (ch) { + struct CONFIG_CONN_HANDLE* next = ch->next; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[INSTANCE_DESTROY] closing config handle for node=0x%016llx", (unsigned long long)ch->node_id); + node_conn_direct_close(ch->handle); + u_free(ch); + ch = next; + } + instance->config_conn_handles = NULL; + // Cleanup ETCP sockets and connections FIRST (before destroying uasync) DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Cleaning up ETCP sockets and connections"); struct ETCP_SOCKET* sock = instance->etcp_sockets; @@ -857,67 +869,52 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc } os = next; } - // Clients [client:] and links (link=) + // Clients [client:] — close old config handles, rebuild from new config + struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles; + while (ch) { struct CONFIG_CONN_HANDLE* next = ch->next; node_conn_direct_close(ch->handle); u_free(ch); ch = next; } + instance->config_conn_handles = NULL; + for (struct CFG_CLIENT *nc = new_config->clients; nc; nc = nc->next) { - struct ETCP_CONN *conn = NULL; - { - struct ll_entry* entry = instance->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (strcmp(nc->name, ce->conn->name) == 0) { conn = ce->conn; break; } - entry = entry->next; - } - } - if (!conn) { - conn = etcp_connection_create(instance, nc->name); - if (conn) { - sc_set_peer_public_key(&conn->crypto_ctx, (const uint8_t*)nc->peer_public_key_hex, 1); - if (random_bytes((uint8_t*)&conn->session_id, sizeof(conn->session_id)) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "Failed to generate session_id for client %s", nc->name); - conn->session_id = 0; - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Added new client %s session_id=%08x", nc->name, conn->session_id); - } - } - // links + if (strlen(nc->peer_public_key_hex) == 0) { DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "client %s no peer pubkey", nc->name); continue; } + uint8_t pubkey_bin[SC_PUBKEY_SIZE]; + if (sc_hex_to_binary(nc->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid pubkey hex client %s", nc->name); continue; } + uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin); + + struct NODE_CONN_DIRECT* handle = NULL; + struct ETCP_CONN* conn = NULL; for (struct CFG_CLIENT_LINK *nl = nc->links; nl; nl = nl->next) { - struct ETCP_LINK *link = NULL; - for (struct ETCP_LINK *l = conn->links; l; l = l->next) { - if (local_sockaddr_equal(&nl->remote_addr, &l->remote_addr)) { - link = l; break; + struct ETCP_SOCKET* sock = NULL; + for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next) + if (nl->local_srv && strcmp(nl->local_srv->name, s->name) == 0) { sock = s; break; } + if (!sock) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s", nl->local_srv ? nl->local_srv->name : "?", nc->name); continue; } + if (!handle) { + struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr; + struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp)); + ni_tmp.node_id = node_id; ni_tmp.node_name = nc->name; + memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE); + if (nl->remote_addr.ss_family == AF_INET) { + struct sockaddr_in* sin = (struct sockaddr_in*)&nl->remote_addr; + v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4); + v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP; + ni_tmp.v4_addrs = &v4_addr; + } else if (nl->remote_addr.ss_family == AF_INET6) { + struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&nl->remote_addr; + v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16); + v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP; + ni_tmp.v6_addrs = &v6_addr; } - } - if (link) { - // in-place for unchanged link (no tear) - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Link for %s unchanged - untouched", nc->name); - continue; - } - struct ETCP_SOCKET *sock = NULL; - for (struct ETCP_SOCKET *s = instance->etcp_sockets; s; s = s->next) { - if (strcmp(nl->server_name, s->name) == 0) { - sock = s; break; - } - } - if (sock) { + int rc = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock); + if (rc == NCD_ERR || !handle) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd reload open failed client %s", nc->name); break; } + conn = node_conn_direct_get_conn(handle); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "reload: client %s ncd %s node=0x%016llx", nc->name, rc == NCD_NEW ? "NEW" : "REUSED", (unsigned long long)node_id); + } else { etcp_link_new(conn, sock, &nl->remote_addr, 0); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Added new link for %s", nc->name); } } - // remove deleted links (safe with next) - struct ETCP_LINK *l = conn->links; - while (l) { - struct ETCP_LINK *next_l = l->next; - int found = 0; - for (struct CFG_CLIENT_LINK *nl = nc->links; nl; nl = nl->next) { - if (local_sockaddr_equal(&nl->remote_addr, &l->remote_addr)) { - found = 1; break; - } - } - if (!found) { - etcp_link_close(l); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Removed deleted link for %s", nc->name); - } - l = next_l; + if (handle) { + ch = u_calloc(1, sizeof(*ch)); + if (ch) { ch->node_id = node_id; strncpy(ch->name, nc->name, MAX_CONN_NAME_LEN - 1); + ch->handle = handle; ch->next = instance->config_conn_handles; instance->config_conn_handles = ch; } } } // Reload networks: clear and repopulate