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.
 
 
 
 
 
 

700 lines
36 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);
/* ═══════════ реестр ═══════════ */
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; }
}
/* ═══════════ поиск узла ═══════════ */
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) ═══════════ */
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;
/* ---- 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;
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (a->port == 0 || !(a->protocol & TOPO_PROTO_UDP)) continue;
int zero = 1; for (int j = 0; j < 4; j++) if (a->addr[j] != 0) { zero = 0; break; }
if (zero) 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));
struct ETCP_SOCKET* use_sock = socks[rr++ % sock_count];
{
struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa);
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: %d.%d.%d.%d:%d already used by live conn 0x%016llx for node 0x%016llx",
a->addr[0], a->addr[1], a->addr[2], a->addr[3], (int)a->port,
(unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id);
continue;
}
}
}
if (etcp_link_new(conn, use_sock, &sa, 0)) link_count++;
}
}
}
/* ---- 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;
for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) {
if (a->port == 0 || !(a->protocol & TOPO_PROTO_UDP)) continue;
int zero = 1; for (int j = 0; j < 16; j++) if (a->addr[j] != 0) { zero = 0; break; }
if (zero) continue;
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);
struct ETCP_SOCKET* use_sock = socks[rr++ % sock_count];
if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = use_sock->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
{
struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa);
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 v6 link node=0x%016llx(old) -> 0x%016llx(new)",
(unsigned long long)old->peer_node_id, (unsigned long long)conn->peer_node_id);
etcp_link_close(stale);
etcp_cbk_fire(old, ETCP_CBK_EVENT_NODE_CHANGED);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "[ncd] addr conflict: v6 port=%d already used by live conn 0x%016llx for node 0x%016llx",
(int)a->port, (unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id);
continue;
}
}
}
if (etcp_link_new(conn, use_sock, &sa, 0)) link_count++;
}
}
}
return link_count;
}
/* ═══════════ единая диспетчеризация событий ═══════════ */
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) ═══════════ */
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;
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);
}
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;
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);
}
/* ═══════════ таймер подключения ═══════════ */
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 handles=%d",
(unsigned long long)entry->node_id, entry->handle_count);
ncd_event_dispatch(entry, NCD_EVENT_TIMEOUT);
}
/* ═══════════ асинхронная доставка UP ═══════════ */
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 ═══════════ */
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; }
}
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 ═══════════ */
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 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:
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);
}
/* ═══════════ инициализация глобального обработчика ═══════════ */
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 ═══════════ */
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 {
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_WARN(DEBUG_CATEGORY_NCD, "[ncd] no links created for node=0x%016llx", (unsigned long long)node_id);
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;
}
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) {
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 {
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_WARN(DEBUG_CATEGORY_NCD, "[ncd] open_node no links created for node=0x%016llx", (unsigned long long)node_id);
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;
}
}
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--;
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->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) {
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);
}
struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h) {
if (!h || !h->entry) return NULL;
return h->entry->conn;
}