You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

1108 lines
60 KiB

/*
* node_conn_direct.c — handle-based ETCP connection layer
*
* Тонкая прослойка над ETCP: управляет прямыми подключениями по node_id.
* Несколько handle'ов могут разделять одно ETCP_CONN.
* Conn закрывается при закрытии последнего handle (с протоколом CLOSE/KEEP_ALIVE).
*/
#include "node_conn_direct.h"
#include "etcp_api.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "secure_channel.h"
#include "utun_instance.h"
#include "topo_node.h"
#include "topo_node_sqlite.h"
#include "topo_group.h"
#include "config_parser.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "../lib/socket_compat.h"
#include "../lib/platform_compat.h"
#include <string.h>
#define DEBUG_CATEGORY_NCD DEBUG_CATEGORY_ETCP
/* ═══════════ локальные константы ═══════════ */
#define NCD_FIN_WAIT_TIMEOUT_TB 50000 /* 5 сек в 0.1ms */
/* ═══════════ внутренние структуры ═══════════ */
struct ncd_entry {
uint64_t node_id;
struct ETCP_CONN* conn;
struct UASYNC* ua;
int handle_count;
struct NODE_CONN_DIRECT* handles; /* связный список всех handle'ов */
void* connect_timer; /* однократный таймер первого подъёма */
void* fin_wait_timer; /* таймер ожидания ответа на CLOSE */
uint8_t up : 1; /* текущий статус: 1=есть живой линк */
uint8_t timed_out : 1; /* connect_timer уже сработал */
struct ncd_entry* next; /* цепочка статического реестра */
};
struct NODE_CONN_DIRECT {
struct ncd_entry* entry;
ncd_callback cb;
void* cb_arg;
struct NODE_CONN_DIRECT* next;
};
/* ─── 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);
/* ═══════════ реестр ═══════════ */
/*
* Реестр ncd_entry по node_id — собственный связный список (не inst->connections).
* ncd_entry хранит состояние, которого нет в ETCP_CONN: handles, таймеры, fin_wait.
* Нужен чтобы при повторном open не создавать дублирующий conn для того же node_id.
*/
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 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 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; }
}
/* ═══════════ поиск узла ═══════════ */
/*
* Загружает информацию об узле (адреса, pubkey) для создания линков.
* Сначала ищет в памяти (node_registry — туда попадают узлы из BGP/topo),
* если нет в памяти — подгружает из SQLite и помещает в реестр.
* Возвращает владеющую ссылку (ref++), вызывающий обязан сделать topo_node_registry_unref.
*/
static struct TOPO_NODE* ncd_lookup_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
if (inst->topo_groups && inst->topo_groups->node_registry) {
struct ll_entry* e = queue_find_data_by_index(inst->topo_groups->node_registry, (const uint8_t*)&node_id);
if (e) {
struct TOPO_NODE* node;
memcpy(&node, e->data + 8, sizeof(node));
if (node) { topo_node_registry_ref(inst->topo_groups, node_id); return node; }
}
}
if (inst->topo_sqlite_db && inst->topo_groups) {
struct TOPO_NODE* ni = topo_node_sqlite_node_load(inst->topo_sqlite_db, inst->topo_groups, node_id);
if (ni) return topo_node_registry_store(inst->topo_groups, ni);
}
return NULL;
}
/* ═══════════ создание линков (round‑robin, все сокеты кроме PRIVATE) ═══════════ */
/*
* Создаёт по одному ETCP_LINK на каждый адрес пира (IPv4 и IPv6).
*
* Все адреса пира распределяются round-robin по локальным сокетам.
* Если specific_sock задан — использует только его. PRIVATE-сокеты пропускаются.
*
* Особые случаи:
* - Нулевые адреса и нулевые порты пропускаются
* - Не-UDP протоколы пропускаются
* - IPv6 link-local: автоматически выставляется scope_id по netif_index сокета
* - Stale-линки: если на том же addr:port висит старый conn с другим pubkey
* и старый conn ещё не поднялся — stale-линк вытесняется (узел пересоздался
* с новым ключом, старый conn больше не нужен)
*/
/* Создаёт один UDP-линк на use_sock → sa, с вытеснением stale-линка и дедупом. */
static int ncd_add_udp_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* use_sock, struct sockaddr_storage* sa) {
struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, sa, 0);
if (stale && stale->etcp == conn) return 0;
if (stale && stale->etcp != conn
&& memcmp(conn->crypto_ctx.peer_public_key, stale->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE))
{
if (!stale->etcp->links_up || !stale->etcp->initialized) {
struct ETCP_CONN* old = stale->etcp;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[ncd] evict stale link node=0x%016llx(old pk=%016llx) for new node=0x%016llx(pk=%016llx)",
(unsigned long long)old->peer_node_id, *(const uint64_t*)old->crypto_ctx.peer_public_key,
(unsigned long long)conn->peer_node_id, *(const uint64_t*)conn->crypto_ctx.peer_public_key);
etcp_link_close(stale);
etcp_cbk_fire(old, ETCP_CBK_EVENT_NODE_CHANGED);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "[ncd] addr conflict: port already used by live conn 0x%016llx for node 0x%016llx",
(unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id);
return 0;
}
}
return etcp_link_new(conn, use_sock, sa, 0) ? 1 : 0;
}
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, v4_skip = 0, v6_skip = 0, v4_addrs = 0, v6_addrs = 0;
uint64_t nid = entry->node_id;
/* ---- IPv4 ---- */
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
v4_addrs++;
if (a->port == 0 || !(a->protocol & (TOPO_PROTO_UDP | TOPO_PROTO_TCP))) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[ncd] v4 skip proto=%02x port=%d node=0x%016llx", a->protocol, (int)a->port, (unsigned long long)nid); v4_skip++; continue; }
int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j] != 0) { zero = 0; break; }
if (zero) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[ncd] v4 skip zero-addr node=0x%016llx", (unsigned long long)nid); continue; }
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
if (a->protocol & TOPO_PROTO_TCP) {
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) {
if (!s->is_tcp || s->local_addr.ss_family != AF_INET || s->type == CFG_SERVER_TYPE_PRIVATE) continue;
struct ETCP_LINK* tlink = etcp_link_new(conn, s, &sa, 0);
if (tlink) { tlink->is_tcp = 1; etcp_tcp_link_start_connect(tlink, &sa, a->port); link_count++; }
}
}
if (a->protocol & TOPO_PROTO_UDP) {
if (specific_sock && !specific_sock->is_tcp && specific_sock->local_addr.ss_family == AF_INET) {
link_count += ncd_add_udp_link(conn, specific_sock, &sa);
} else if (!specific_sock) {
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) {
if (s->is_tcp || s->local_addr.ss_family != AF_INET || s->type == CFG_SERVER_TYPE_PRIVATE) continue;
link_count += ncd_add_udp_link(conn, s, &sa);
}
}
}
}
/* ---- IPv6 ---- */
for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) {
v6_addrs++;
if (a->port == 0 || !(a->protocol & (TOPO_PROTO_UDP | TOPO_PROTO_TCP))) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[ncd] v6 skip proto=%02x port=%d node=0x%016llx", a->protocol, (int)a->port, (unsigned long long)nid); v6_skip++; continue; }
int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j] != 0) { zero = 0; break; }
if (zero) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[ncd] v6 skip zero-addr node=0x%016llx", (unsigned long long)nid); continue; }
int is_ll = (a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80);
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port);
if (a->protocol & TOPO_PROTO_TCP) {
if (is_ll) { struct ETCP_SOCKET* sv = inst->etcp_sockets;
while (sv) { if (sv->local_addr.ss_family == AF_INET6 && sv->netif_index) { sin6.sin6_scope_id = sv->netif_index; break; } sv = sv->next; } }
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) {
if (!s->is_tcp || s->local_addr.ss_family != AF_INET6 || s->type == CFG_SERVER_TYPE_PRIVATE) continue;
struct ETCP_LINK* tlink = etcp_link_new(conn, s, &sa, 0);
if (tlink) { tlink->is_tcp = 1; etcp_tcp_link_start_connect(tlink, &sa, a->port); link_count++; }
}
}
if (a->protocol & TOPO_PROTO_UDP) {
if (specific_sock && !specific_sock->is_tcp && specific_sock->local_addr.ss_family == AF_INET6) {
if (is_ll) sin6.sin6_scope_id = specific_sock->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
link_count += ncd_add_udp_link(conn, specific_sock, &sa);
} else if (!specific_sock) {
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) {
if (s->is_tcp || s->local_addr.ss_family != AF_INET6 || s->type == CFG_SERVER_TYPE_PRIVATE) continue;
if (is_ll) sin6.sin6_scope_id = s->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
link_count += ncd_add_udp_link(conn, s, &sa);
}
}
}
}
if (link_count == 0) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] create_links: node=0x%016llx v4=%d/skp=%d v6=%d/skp=%d → %d links (UDP+TCP)",
(unsigned long long)nid, v4_addrs, v4_skip, v6_addrs, v6_skip, link_count);
}
return link_count;
}
/* ═══════════ единая диспетчеризация событий ═══════════ */
/*
* Центральный диспетчер: принимает событие (UP/DOWN/TIMEOUT), проверяет что оно
* действительно меняет состояние (повторный UP когда уже up — игнорируется),
* обновляет entry->up/timed_out и рассылает callback всем handle'ам.
*
* Именно через эту функцию все потребители узнают об изменении состояния соединения.
* Один вызов — один проход по всем handle'ам.
*/
static void ncd_event_dispatch(struct ncd_entry* entry, enum ncd_event event) {
if (!entry) return;
switch (event) {
case NCD_EVENT_UP:
if (entry->up) return;
if (entry->connect_timer) { uasync_cancel_timeout(entry->ua, entry->connect_timer); entry->connect_timer = NULL; }
entry->up = 1;
break;
case NCD_EVENT_DOWN:
if (!entry->up) return;
entry->up = 0;
break;
case NCD_EVENT_TIMEOUT:
if (entry->timed_out) return;
entry->timed_out = 1;
entry->connect_timer = NULL;
break;
}
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] event=%s node=0x%016llx handles=%d",
event == NCD_EVENT_UP ? "UP" : event == NCD_EVENT_DOWN ? "DOWN" : "TIMEOUT",
(unsigned long long)entry->node_id, entry->handle_count);
struct NODE_CONN_DIRECT* h = entry->handles;
while (h) { struct NODE_CONN_DIRECT* next = h->next; if (h->cb) h->cb(h, event, h->cb_arg); h = next; }
}
/* ═══════════ ETCP коллбэки (прокидывают в ncd_event_dispatch) ═══════════ */
/*
* ETCP сообщил что соединение работает: INIT — handshake завершён,
* UP — линки восстановились после DOWN. Транслируем в NCD_EVENT_UP.
* Не транслируем если уже сработал таймаут подключения (timed_out)
* или conn в состоянии fin_wait (ждём закрытия).
* Если INIT пришёл а линков нет — UP не доставляем: ncd_up_cb сам
* доставит когда линки поднимутся. Так гарантируем что пользователь
* получает UP только когда conn реально готов отправлять данные.
*/
static void ncd_init_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)event;
struct ncd_entry* entry = (struct ncd_entry*)arg;
if (!entry || entry->timed_out || conn->fin_wait) return;
if (conn->links_up > 0) ncd_event_dispatch(entry, NCD_EVENT_UP);
}
static void ncd_up_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)event;
struct ncd_entry* entry = (struct ncd_entry*)arg;
if (!entry || conn->fin_wait) return;
ncd_event_dispatch(entry, NCD_EVENT_UP);
}
/*
* Отложенное (uasync_call_soon) уничтожение ncd_entry вместе с conn.
* Нельзя вызывать напрямую из ETCP-коллбэка: conn может использоваться
* после возврата из коллбэка. Поэтому очистка откладывается на следующий цикл событий.
* Снимает все ETCP-коллбэки, удаляет из реестра, закрывает conn, освобождает память.
*/
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);
}
/*
* Отложенное закрытие conn без ncd_entry.
* Используется когда пришёл CLOSE от пира, а ncd_entry для этого пира уже нет
* (закрыли раньше — например другой модуль уже удалил все handle'ы).
*/
static void ncd_deferred_close_conn(void* arg) {
struct ETCP_CONN* conn = (struct ETCP_CONN*)arg;
if (conn) etcp_connection_close(conn);
}
/*
* ETCP сообщил что все линки упали.
* Если conn в fin_wait и handle'ов нет — peer подтвердил закрытие,
* делаем немедленную очистку. Иначе транслируем NCD_EVENT_DOWN потребителям.
*/
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;
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; }
uasync_call_soon(entry->ua, entry, ncd_deferred_close);
return;
}
ncd_event_dispatch(entry, NCD_EVENT_DOWN);
}
/* ═══════════ таймер подключения ═══════════ */
/*
* Таймер первого подключения истёк — соединение не установилось за etcp_connect_timeout_tb.
* Транслирует NCD_EVENT_TIMEOUT всем handle'ам. После этого init/up коллбэки
* больше не транслируются (entry->timed_out = 1).
*/
static void ncd_connect_timeout_cb(void* arg) {
struct ncd_entry* entry = (struct ncd_entry*)arg;
if (!entry || entry->timed_out) return;
entry->connect_timer = NULL;
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] connect TIMEOUT: node=0x%016llx now_tb=%llu handles=%d",
(unsigned long long)entry->node_id, (unsigned long long)get_time_tb(), entry->handle_count);
ncd_event_dispatch(entry, NCD_EVENT_TIMEOUT);
}
/* ═══════════ асинхронная доставка UP ═══════════ */
/*
* Доставляет NCD_EVENT_UP одному конкретному handle'у (не всем).
* Используется когда handle добавляется к уже работающему conn: доставка отложенная
* (uasync_call_soon), чтобы вызывающий успел сохранить указатель handle
* до вызова callback'а. Синхронная доставка привела бы к гонке.
*/
static void ncd_deliver_up_cb(void* arg) {
struct NODE_CONN_DIRECT* h = (struct NODE_CONN_DIRECT*)arg;
if (h->cb) h->cb(h, NCD_EVENT_UP, h->cb_arg);
}
/* ═══════════ FIN_WAIT ═══════════ */
/*
* fin_wait — состояние ожидания подтверждения закрытия от пира.
* Вызывающая сторона: мы отправили CLOSE и ждём ответа (KEEP_ALIVE или DOWN).
* Максимальное время ожидания — NCD_FIN_WAIT_TIMEOUT_TB (5 сек), после чего
* соединение форсированно закрывается.
*/
/*
* ETCP отменяет наш fin_wait — peer переподключился и прислал новый INIT
* (или другая причина отмены на уровне ETCP). Отменяем локальный таймер,
* соединение продолжает работать.
*/
static void ncd_fin_wait_cancelled(struct ETCP_CONN* conn, void* arg) {
struct ncd_entry* entry = (struct ncd_entry*)arg;
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] fin_wait cancelled node=0x%016llx", (unsigned long long)entry->node_id);
if (entry->fin_wait_timer) { uasync_cancel_timeout(entry->ua, entry->fin_wait_timer); entry->fin_wait_timer = NULL; }
}
/*
* Таймаут fin_wait: peer не ответил на CLOSE за NCD_FIN_WAIT_TIMEOUT_TB (5 сек).
* Форсированно закрываем conn и освобождаем entry — соединение разорвано
* без подтверждения от пира.
*/
static void ncd_fin_wait_timeout_cb(void* arg) {
struct ncd_entry* entry = (struct ncd_entry*)arg;
if (!entry || !entry->conn) return;
entry->fin_wait_timer = NULL;
if (!entry->conn->fin_wait) return;
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] fin_wait timeout, force close node=0x%016llx", (unsigned long long)entry->node_id);
entry->conn->fin_wait = 0; entry->conn->fin_wait_clear_cb = NULL; entry->conn->fin_wait_clear_arg = NULL;
etcp_conn_remove_cbk(entry->conn, ncd_init_cb, entry);
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->conn->instance, entry);
u_free(entry);
}
/* ═══════════ приём CLOSE / KEEP_ALIVE ═══════════ */
/*
* Обрабатывает управляющие сообщения от пира по протоколу graceful shutdown.
*
* CLOSE (пир хочет закрыть conn — у него закончились handle'ы):
* - Если у нас ещё есть handle'ы — отказываем: шлём KEEP_ALIVE, conn живёт.
* - Если handle'ов нет — соглашаемся, закрываем conn.
*
* KEEP_ALIVE (пир просит не закрывать conn — у него ещё есть handle'ы):
* - Если мы в fin_wait — отменяем его, conn продолжает работать.
* - Если у нас handle'ов нет — перезапускаем fin_wait (ждём ещё 5 сек).
*/
static void ncd_recv_control_handler(struct ETCP_CONN* conn, struct ll_entry* e) {
if (!conn || !e || e->len < NCD_CONTROL_MSG_SIZE) {
if (e) { queue_dgram_free(e); queue_entry_free(e); }
return;
}
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(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, 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) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] refusing CLOSE, still have %d handles for 0x%016llx",
entry->handle_count, (unsigned long long)sender_id);
{ struct ncd_control_msg resp; resp.cmd = ETCP_RT_ID_NCD_CONTROL; resp.subcmd = NCD_SUBCMD_KEEP_ALIVE;
resp.node_id = conn->instance->node_id;
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = u_malloc(sizeof(resp)); memcpy(qe->dgram, &resp, sizeof(resp)); qe->len = sizeof(resp);
etcp_send(conn, qe); }
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] deferred close entry for CLOSE from 0x%016llx", (unsigned long long)sender_id);
if (entry->conn && entry->conn->conn_queue && entry->conn->conn_queue_entry) {
queue_remove_data(entry->conn->conn_queue, entry->conn->conn_queue_entry);
queue_entry_free(entry->conn->conn_queue_entry);
entry->conn->conn_queue_entry = NULL; entry->conn->conn_queue = NULL;
}
uasync_call_soon(entry->ua, entry, ncd_deferred_close);
}
break;
case NCD_SUBCMD_KEEP_ALIVE:
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] recv KEEP_ALIVE from 0x%016llx", (unsigned long long)sender_id);
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);
}
/* ═══════════ инициализация глобального обработчика ═══════════ */
/*
* Однократная привязка обработчика NCD control-сообщений к ETCP (etcp_bind).
* Вызывается при первом node_conn_direct_open, повторные вызовы безвредны
* (защита через inst->ncd_control_bound).
*/
static void ncd_init_control_binding(struct UTUN_INSTANCE* inst) {
if (inst->ncd_control_bound) return;
etcp_bind(inst, ETCP_RT_ID_NCD_CONTROL, ncd_recv_control_handler);
inst->ncd_control_bound = 1;
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] bound ETCP_RT_ID_NCD_CONTROL handler");
}
/* ═══════════ API ═══════════ */
/*
* Открыть handle для связи с удалённым узлом.
*
* Три сценария (прозрачно для вызывающего):
* 1. Узел уже в NCD-реестре — другой модуль уже открыл соединение.
* Добавляем ещё один handle. Если conn работает — сразу шлём UP.
* 2. ETCP-соединение существует (входящее / etcp_connect), но NCD о нём не знает.
* Оборачиваем в ncd_entry, подписываемся на события.
* 3. Ничего нет — загружаем адреса/pupkey узла, создаём новый ETCP_CONN,
* инициализируем шифрование, создаём линки, запускаем таймер подключения.
*
* Возвращает NCD_NEW (новое), NCD_REUSED (переиспользовано) или NCD_ERR.
* specific_sock=NULL — авто-подбор всех локальных сокетов.
*/
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 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(inst, node_id);
if (entry) {
/* снять fin_wait если был */
if (entry->conn && entry->conn->fin_wait) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED clearing fin_wait node=0x%016llx", (unsigned long long)node_id);
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; }
}
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); return NCD_ERR; }
h->entry = entry; h->cb = cb; h->cb_arg = cb_arg;
h->next = entry->handles; entry->handles = h;
entry->handle_count++;
*out_handle = h;
if (entry->conn && entry->conn->state == 1 && entry->up) {
uasync_call_soon(inst->ua, h, ncd_deliver_up_cb);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED (ready) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
} else {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED (pending) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
}
return NCD_REUSED;
}
/* 2. Ищем conn через instance_find_conn (входящее / созданное etcp_connect) */
struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (conn && conn->state == 1) {
/* снять fin_wait если был */
if (conn->fin_wait) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED clearing fin_wait (existing conn) node=0x%016llx", (unsigned long long)node_id);
conn->fin_wait = 0; conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL;
}
entry = u_calloc(1, sizeof(*entry));
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(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(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;
*out_handle = h;
etcp_conn_add_cbk(conn, ncd_init_cb, entry, ETCP_CBK_EVENT_INIT);
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);
if (conn->state == 1 && entry->up) {
uasync_call_soon(inst->ua, h, ncd_deliver_up_cb);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED new-entry (ready) node=0x%016llx conn=%p", (unsigned long long)node_id, (void*)conn);
} else {
{ struct TOPO_NODE* ni = ncd_lookup_node(inst, node_id);
if (ni) { ncd_create_links(entry, ni, specific_sock); topo_node_registry_unref(inst->topo_groups, ni->node_id); }
}
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "[ncd] connect timer SET: node=0x%016llx value_tb=%u now_tb=%llu path=REUSED_entry_pending",
(unsigned long long)node_id, inst->etcp_connect_timeout_tb, (unsigned long long)get_time_tb());
entry->connect_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, entry, ncd_connect_timeout_cb, "ncd_connect");
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED new-entry (pending) node=0x%016llx conn=%p", (unsigned long long)node_id, (void*)conn);
}
return NCD_REUSED;
}
/* 3. Новое подключение */
struct TOPO_NODE* ni = ncd_lookup_node(inst, node_id);
if (!ni) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] node not found id=0x%016llx", (unsigned long long)node_id);
return NCD_ERR;
}
{ 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);
return NCD_ERR;
}
conn->peer_node_id = node_id;
etcp_update_log_name(conn);
if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] sc_init_ctx failed node=0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); topo_node_registry_unref(inst->topo_groups, ni->node_id); return NCD_ERR;
}
if (sc_set_peer_public_key(&conn->crypto_ctx, ni->public_key, 0) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] sc_set_peer_public_key failed node=0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); topo_node_registry_unref(inst->topo_groups, ni->node_id); return NCD_ERR;
}
/* Подключаем conn в inst->connections с индексом по node_id */
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
{ struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (qe) { struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = node_id; ce->conn = conn;
conn->conn_queue_entry = qe; conn->conn_queue = inst->connections;
queue_data_put_with_index(inst->connections, qe); }
}
entry = u_calloc(1, sizeof(*entry));
if (!entry) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] alloc entry failed node=0x%016llx", (unsigned long long)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(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(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;
*out_handle = h;
etcp_conn_add_cbk(conn, ncd_init_cb, entry, ETCP_CBK_EVENT_INIT);
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, specific_sock);
if (link_count == 0)
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] no links created (UDP only, TCP added by conn_mgr) node=0x%016llx", (unsigned long long)node_id);
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "[ncd] connect timer SET: node=0x%016llx value_tb=%u now_tb=%llu path=NEW_node_conn_direct_open",
(unsigned long long)node_id, inst->etcp_connect_timeout_tb, (unsigned long long)get_time_tb());
entry->connect_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, entry, ncd_connect_timeout_cb, "ncd_connect");
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open NEW node=0x%016llx conn=%p links=%d handles=%d",
(unsigned long long)node_id, (void*)conn, link_count, entry->handle_count);
topo_node_registry_unref(inst->topo_groups, ni->node_id);
return NCD_NEW;
}
/*
* То же что open, но TOPO_NODE (адреса + pubkey) уже загружен вызывающим.
* Экономит поиск узла (ncd_lookup_node) — полезно когда адреса/pupkey
* известны заранее (например из BGP-анонса). ni не владеем — можно
* передать временную структуру, копия не делается.
*/
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 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(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);
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; }
}
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); return NCD_ERR; }
h->entry = entry; h->cb = cb; h->cb_arg = cb_arg;
h->next = entry->handles; entry->handles = h;
entry->handle_count++;
*out_handle = h;
if (entry->conn && entry->conn->state == 1 && entry->up) {
uasync_call_soon(inst->ua, h, ncd_deliver_up_cb);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED (ready) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
} else {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED (pending) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
}
return NCD_REUSED;
}}
/* 2. Ищем conn через instance_find_conn */
{ struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (conn && conn->state != 2) {
if (conn->fin_wait) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED clearing fin_wait (existing conn) node=0x%016llx", (unsigned long long)node_id);
conn->fin_wait = 0; conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL;
}
struct ncd_entry* entry = u_calloc(1, sizeof(*entry));
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(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(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;
*out_handle = h;
etcp_conn_add_cbk(conn, ncd_init_cb, entry, ETCP_CBK_EVENT_INIT);
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);
if (conn->state == 1 && entry->up) {
uasync_call_soon(inst->ua, h, ncd_deliver_up_cb);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED new-entry (ready) node=0x%016llx conn=%p", (unsigned long long)node_id, (void*)conn);
} else {
ncd_create_links(entry, ni, specific_sock);
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "[ncd] connect timer SET: node=0x%016llx value_tb=%u now_tb=%llu path=REUSED_conn_pending",
(unsigned long long)node_id, inst->etcp_connect_timeout_tb, (unsigned long long)get_time_tb());
entry->connect_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, entry, ncd_connect_timeout_cb, "ncd_connect_node");
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED new-entry (pending) node=0x%016llx conn=%p", (unsigned long long)node_id, (void*)conn);
}
return NCD_REUSED;
}}
/* 3. Новое подключение — используем переданный ni (временный, не владеем) */
{ 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;
}
conn->peer_node_id = node_id;
etcp_update_log_name(conn);
if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] open_node sc_init_ctx failed node=0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); return NCD_ERR;
}
if (sc_set_peer_public_key(&conn->crypto_ctx, ni->public_key, 0) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] open_node sc_set_peer_public_key failed node=0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); return NCD_ERR;
}
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
{ struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (qe) { struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = node_id; ce->conn = conn;
conn->conn_queue_entry = qe; conn->conn_queue = inst->connections;
queue_data_put_with_index(inst->connections, qe); }
}
struct ncd_entry* entry = u_calloc(1, sizeof(*entry));
if (!entry) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] open_node alloc entry failed node=0x%016llx", (unsigned long long)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(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(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;
*out_handle = h;
etcp_conn_add_cbk(conn, ncd_init_cb, entry, ETCP_CBK_EVENT_INIT);
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, specific_sock);
if (link_count == 0)
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node (UDP only, TCP added by conn_mgr) node=0x%016llx links=%d", (unsigned long long)node_id, link_count);
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "[ncd] connect timer SET: node=0x%016llx value_tb=%u now_tb=%llu path=NEW_node_conn_direct_open_node",
(unsigned long long)node_id, inst->etcp_connect_timeout_tb, (unsigned long long)get_time_tb());
entry->connect_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, entry, ncd_connect_timeout_cb, "ncd_connect_node");
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node NEW node=0x%016llx conn=%p links=%d handles=%d",
(unsigned long long)node_id, (void*)conn, link_count, entry->handle_count);
return NCD_NEW;
}
}
/*
* Закрыть handle (graceful shutdown).
*
* Если остались другие handle'ы — просто удаляется из списка, conn продолжает работу.
*
* Если это был последний handle (handle_count стал 0):
* - Запускается протокол graceful shutdown: пиру отправляется CLOSE, conn
* переводится в fin_wait, ставится таймер на 5 сек.
* - Если пир ответит KEEP_ALIVE (у него ещё есть handle'ы) — conn останется жив,
* просто без наших handle'ов.
* - Если пир не ответит — conn закрывается форсированно по таймауту.
*/
void node_conn_direct_close(struct NODE_CONN_DIRECT* h) {
if (!h) return;
struct ncd_entry* entry = h->entry;
if (!entry) { u_free(h); return; }
uint64_t node_id = entry->node_id;
struct ETCP_CONN* conn = entry->conn;
/* Удаляем handle из списка */
struct NODE_CONN_DIRECT** pp = &entry->handles;
while (*pp) { if (*pp == h) { *pp = h->next; break; } pp = &(*pp)->next; }
entry->handle_count--;
if (entry->handle_count < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] FATAL: handle_count=%d < 0 node=0x%016llx",
entry->handle_count, (unsigned long long)node_id);
h->entry = NULL; u_free(h); return;
}
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] close node=0x%016llx remaining=%d", (unsigned long long)node_id, entry->handle_count);
if (entry->handle_count == 0) {
if (entry->handles) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] FATAL: handle_count=0 but handles non-empty node=0x%016llx",
(unsigned long long)node_id);
entry->handles = NULL;
}
h->cb = NULL; h->cb_arg = NULL;
if (entry->connect_timer) { uasync_cancel_timeout(entry->ua, entry->connect_timer); entry->connect_timer = NULL; }
if (conn && conn->state != 2) {
/* устанавливаем fin_wait, отправляем CLOSE, ставим таймер */
conn->fin_wait = 1;
conn->fin_wait_clear_cb = ncd_fin_wait_cancelled;
conn->fin_wait_clear_arg = entry;
{ struct ncd_control_msg msg; msg.cmd = ETCP_RT_ID_NCD_CONTROL; msg.subcmd = NCD_SUBCMD_CLOSE;
msg.node_id = conn->instance->node_id;
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = u_malloc(sizeof(msg)); memcpy(qe->dgram, &msg, sizeof(msg)); qe->len = sizeof(msg);
etcp_send(conn, qe);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] sent CLOSE to 0x%016llx", (unsigned long long)node_id); }
}
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] fin_wait started node=0x%016llx", (unsigned long long)node_id);
} else {
if (conn && conn->state != 2) {
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(conn->instance, entry);
u_free(entry);
} else {
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] close: no conn for entry, leaking entry=%p", entry);
}
}
}
h->entry = NULL;
u_free(h);
}
/*
* Немедленное жёсткое закрытие handle. Никакого CLOSE/KEEP_ALIVE — сразу
* снимает ETCP-коллбэки и (если последний handle) закрывает conn.
*
* Используется при разрушении CM entry чтобы гарантировать что коллбэки
* не вызовутся в уже освобождённую память. После force_close никакие
* NCD-события не будут доставлены.
*/
void node_conn_direct_force_close(struct NODE_CONN_DIRECT* h) {
if (!h) return;
if (!h->entry) { DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] force_close: h=%p h->entry=NULL", h); u_free(h); return; }
struct ncd_entry* entry = h->entry;
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] force_close h=%p entry=%p node=0x%016llx handles=%d",
h, entry, (unsigned long long)entry->node_id, entry->handle_count);
uint64_t node_id = entry->node_id;
struct ETCP_CONN* conn = entry->conn;
struct NODE_CONN_DIRECT** pp = &entry->handles;
while (*pp) { if (*pp == h) { *pp = h->next; break; } pp = &(*pp)->next; }
entry->handle_count--;
if (entry->handle_count < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] FATAL: force_close handle_count=%d < 0 node=0x%016llx",
entry->handle_count, (unsigned long long)node_id);
h->entry = NULL; u_free(h); return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close node=0x%016llx remaining=%d",
(unsigned long long)node_id, entry->handle_count);
if (entry->handle_count == 0) {
if (entry->handles) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] FATAL: force_close handle_count=0 but handles non-empty node=0x%016llx",
(unsigned long long)node_id);
entry->handles = NULL;
}
h->cb = NULL; h->cb_arg = NULL;
if (entry->connect_timer) { uasync_cancel_timeout(entry->ua, entry->connect_timer); entry->connect_timer = NULL; }
if (entry->fin_wait_timer) { uasync_cancel_timeout(entry->ua, entry->fin_wait_timer); entry->fin_wait_timer = NULL; }
if (conn && conn->state != 2 && conn->instance) {
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: cleaning conn=%p", conn);
conn->fin_wait = 0;
conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL;
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing ncd_init_cb");
etcp_conn_remove_cbk(conn, ncd_init_cb, entry);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing ncd_up_cb");
etcp_conn_remove_cbk(conn, ncd_up_cb, entry);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing ncd_down_cb");
etcp_conn_remove_cbk(conn, ncd_down_cb, entry);
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
conn->conn_queue_entry = NULL; conn->conn_queue = NULL;
}
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: scheduling deferred close");
if (conn->state != 2) uasync_call_soon(entry->ua, conn, ncd_deferred_close_conn);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: removing from registry");
ncd_registry_remove(conn->instance, entry);
}
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: freeing entry=%p", entry);
u_free(entry);
}
h->entry = NULL;
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: freeing h=%p", h);
u_free(h);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] force_close: done");
}
/*
* Сменить или сбросить (cb=NULL) callback на уже открытом handle.
* Не влияет на refcounting и состояние conn.
*/
void node_conn_direct_set_callback(struct NODE_CONN_DIRECT* h, ncd_callback cb, void* cb_arg) {
if (!h || !h->entry) return;
h->cb = cb;
h->cb_arg = cb_arg;
}
void node_conn_direct_transfer(struct NODE_CONN_DIRECT* h, ncd_callback new_cb, void* new_cb_arg) {
if (!h || !h->entry) return;
h->cb = new_cb;
h->cb_arg = new_cb_arg;
}
/*
* Прямой доступ к ETCP_CONN из handle.
* Нужен для отправки данных (etcp_send), проверки статуса, и т.д.
*/
struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h) {
if (!h || !h->entry) return NULL;
return h->entry->conn;
}
/* ═══════════ Управление линками на конкретном сокете (для auto_socket) ═══════════ */
/* ncd_add_udp_link_random / ncd_add_tcp_link_random — собирают совместимые адреса пира,
* случайно выбирают один, проверяют что линк на (sock, addr, proto) ещё не существует,
* и создают новый. */
static int ncd_add_udp_link_random(struct ETCP_CONN* conn, struct TOPO_NODE* ni, struct ETCP_SOCKET* sock) {
int family = sock->local_addr.ss_family;
if (family == AF_INET) {
struct TOPO_ADDR4* compat[64]; int count = 0;
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (!a->port || !(a->protocol & TOPO_PROTO_UDP)) continue;
int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j]) { zero = 0; break; }
if (zero) continue;
if (count < 64) compat[count++] = (struct TOPO_ADDR4*)a;
}
if (count == 0) return 0;
uint32_t r; random_bytes((uint8_t*)&r, sizeof(r));
struct TOPO_ADDR4* a = compat[r % count];
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
if (etcp_link_find_by_addr(sock, &sa, 0)) return 0;
if (etcp_link_new(conn, sock, &sa, 0)) return 1;
return 0;
}
if (family == AF_INET6) {
struct TOPO_ADDR6* compat[64]; int count = 0;
for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) {
if (!a->port || !(a->protocol & TOPO_PROTO_UDP)) continue;
int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j]) { zero = 0; break; }
if (zero) continue;
if (count < 64) compat[count++] = (struct TOPO_ADDR6*)a;
}
if (count == 0) return 0;
uint32_t r; random_bytes((uint8_t*)&r, sizeof(r));
struct TOPO_ADDR6* a = compat[r % count];
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port);
if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = sock->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
if (etcp_link_find_by_addr(sock, &sa, 0)) return 0;
if (etcp_link_new(conn, sock, &sa, 0)) return 1;
return 0;
}
return 0;
}
static int ncd_add_tcp_link_random(struct ETCP_CONN* conn, struct TOPO_NODE* ni, struct ETCP_SOCKET* sock) {
int family = sock->local_addr.ss_family;
if (family == AF_INET) {
struct TOPO_ADDR4* compat[64]; int count = 0;
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (!a->port || !(a->protocol & TOPO_PROTO_TCP)) continue;
int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j]) { zero = 0; break; }
if (zero) continue;
if (count < 64) compat[count++] = (struct TOPO_ADDR4*)a;
}
if (count == 0) return 0;
uint32_t r; random_bytes((uint8_t*)&r, sizeof(r));
struct TOPO_ADDR4* a = compat[r % count];
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
if (etcp_link_find_by_addr(sock, &sa, 1)) return 0;
struct ETCP_LINK* tlink = etcp_link_new(conn, sock, &sa, 0);
if (!tlink) return 0;
tlink->is_tcp = 1;
etcp_tcp_link_start_connect(tlink, &sa, a->port);
return 1;
}
if (family == AF_INET6) {
struct TOPO_ADDR6* compat[64]; int count = 0;
for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) {
if (!a->port || !(a->protocol & TOPO_PROTO_TCP)) continue;
int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j]) { zero = 0; break; }
if (zero) continue;
if (count < 64) compat[count++] = (struct TOPO_ADDR6*)a;
}
if (count == 0) return 0;
uint32_t r; random_bytes((uint8_t*)&r, sizeof(r));
struct TOPO_ADDR6* a = compat[r % count];
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port);
if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = sock->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
if (etcp_link_find_by_addr(sock, &sa, 1)) return 0;
struct ETCP_LINK* tlink = etcp_link_new(conn, sock, &sa, 0);
if (!tlink) return 0;
tlink->is_tcp = 1;
etcp_tcp_link_start_connect(tlink, &sa, a->port);
return 1;
}
return 0;
}
int ncd_add_socket_links(struct UTUN_INSTANCE* inst, struct ETCP_SOCKET* sock) {
if (!inst || !sock) return 0;
int added = 0, total_conns = 0, skipped_closed = 0, skipped_no_addr = 0;
struct ll_entry* qe = inst->connections ? inst->connections->head : NULL;
while (qe) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
struct ETCP_CONN* conn = ce->conn;
qe = qe->next;
if (!conn) continue;
total_conns++;
if (conn->state == 2) { skipped_closed++; continue; }
struct TOPO_NODE* ni = ncd_lookup_node(inst, ce->peer_node_id);
if (!ni) { skipped_no_addr++; continue; }
int n;
if (sock->is_tcp) n = ncd_add_tcp_link_random(conn, ni, sock);
else n = ncd_add_udp_link_random(conn, ni, sock);
if (n > 0) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] add_socket_links sock=%s fd=%d conn=[%s] node=0x%016llx addr=%s is_tcp=%d",
sock->name, (int)sock->fd, conn->log_name,
(unsigned long long)ce->peer_node_id,
sockaddr_storage_to_str(&sock->local_addr).str, sock->is_tcp);
added += n;
}
topo_node_registry_unref(inst->topo_groups, ni->node_id);
}
DEBUG_INFO(DEBUG_CATEGORY_NCD,
"[ncd] add_socket_links done: sock=%s is_tcp=%d total_conns=%d skipped(cl=%d,no_addr=%d) added=%d",
sock->name, sock->is_tcp, total_conns, skipped_closed, skipped_no_addr, added);
return added;
}
int ncd_remove_socket_links(struct UTUN_INSTANCE* inst, struct ETCP_SOCKET* sock) {
if (!inst || !sock) return 0;
int removed = 0;
struct ncd_entry* entry = (struct ncd_entry*)inst->ncd_registry;
while (entry) {
struct ETCP_CONN* conn = entry->conn;
if (conn) {
struct ETCP_LINK* link = conn->links;
while (link) {
struct ETCP_LINK* next = link->next;
if (link->conn == sock) {
etcp_link_close(link);
removed++;
}
link = next;
}
}
entry = entry->next;
}
if (removed > 0) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] remove_socket_links socket=%s closed=%d", sock->name, removed);
}
return removed;
}