Browse Source

node_conn_direct: add open_node variant, refactor internals

topo_upd
evgeny 2 months ago
parent
commit
eb2a8f6975
  1. 323
      src/transport_layer/node_conn_direct.c
  2. 16
      src/transport_layer/node_conn_direct.h
  3. 135
      tests/test_node_conn_direct.c

323
src/transport_layer/node_conn_direct.c

@ -1,18 +1,15 @@
/*
* node_conn_direct.c — handle-based ETCP connection layer
*
* Тонкая прослойка над ETCP: управляет подключениями по node_id.
* Тонкая прослойка над ETCP: управляет прямыми подключениями по node_id.
* Несколько handle'ов могут разделять одно ETCP_CONN.
* Conn закрывается при закрытии последнего handle (с протоколом CLOSE/KEEP_ALIVE).
* Поддерживает REVERSE-подключение через etcp_router если задан group_id.
*/
#include "node_conn_direct.h"
#include "etcp_api.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "etcp_router.h"
#include "conn_mgr.h"
#include "secure_channel.h"
#include "utun_instance.h"
#include "topo_node.h"
@ -31,18 +28,12 @@
/* ═══════════ локальные константы ═══════════ */
#define CM_V6_LINK_LOCAL 0
#define CM_V6_LOCAL 1
#define CM_V6_DIRECT 2
#define CM_V6_OTHER 3
#define NCD_FIN_WAIT_TIMEOUT_TB 50000 /* 5 сек в 0.1ms */
/* ═══════════ внутренние структуры ═══════════ */
struct ncd_entry {
uint64_t node_id;
uint64_t group_id; /* для REVERSE (0 = без REVERSE) */
struct ETCP_CONN* conn;
struct UASYNC* ua;
int handle_count;
@ -51,7 +42,6 @@ struct ncd_entry {
void* connect_timer; /* однократный таймер первого подъёма */
void* fin_wait_timer; /* таймер ожидания ответа на CLOSE */
void* reverse_poll_timer; /* таймер polling входящего conn при REVERSE */
uint8_t up : 1; /* текущий статус: 1=есть живой линк */
uint8_t timed_out : 1; /* connect_timer уже сработал */
@ -158,99 +148,6 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni) {
return link_count;
}
/* ═══════════ REVERSE ═══════════ */
static uint8_t ncd_classify_v6_addr(const uint8_t addr[16]) {
if (addr[0] == 0xfe && (addr[1] & 0xc0) == 0x80) return CM_V6_LINK_LOCAL;
if (addr[0] == 0xfc || addr[0] == 0xfd) return CM_V6_LOCAL;
if (addr[0] == 0xff) return CM_V6_OTHER;
{ uint8_t zero[16] = {0}; if (memcmp(addr, zero, 16) == 0) return CM_V6_OTHER; }
{ uint8_t lb[16] = {0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,1}; if (memcmp(addr, lb, 16) == 0) return CM_V6_OTHER; }
return CM_V6_DIRECT;
}
static int ncd_has_direct_ip(const struct TOPO_NODE* node) {
if (!node) return 0;
for (const struct TOPO_ADDR4* a = node->v4_addrs; a; a = a->next) {
if (a->type == TOPO_ADDR_NAT || a->type == TOPO_ADDR_INTERFACE) {
for (const struct TOPO_SOCKMETA4* m = node->v4_sock_meta; m; m = m->next) {
if (m->id == a->socket_id && (m->config_type == CFG_SERVER_TYPE_PUBLIC || m->config_type == CFG_SERVER_TYPE_UNKNOWN
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT))
return 1;
}
}
}
for (const struct TOPO_ADDR6* a6 = node->v6_addrs; a6; a6 = a6->next) {
if ((a6->type != TOPO_ADDR_NAT && a6->type != TOPO_ADDR_INTERFACE) || ncd_classify_v6_addr(a6->addr) != CM_V6_DIRECT) continue;
for (const struct TOPO_SOCKMETA6* m = node->v6_sock_meta; m; m = m->next)
if (m->id == a6->socket_id && (m->config_type == CFG_SERVER_TYPE_PUBLIC || m->config_type == CFG_SERVER_TYPE_UNKNOWN
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT))
return 1;
}
return 0;
}
static void ncd_try_reverse(struct ncd_entry* entry, struct TOPO_NODE* local_node, struct TOPO_NODE* target_node) {
if (!local_node || !target_node) return;
int has_our = ncd_has_direct_ip(local_node);
int has_target = ncd_has_direct_ip(target_node);
if (!has_our || has_target) return;
uint8_t direct_count = 0;
struct { uint8_t type; uint8_t ip[4]; uint16_t port; uint8_t socket_id; } out_addrs[8];
for (const struct TOPO_ADDR4* a = local_node->v4_addrs; a && direct_count < 8; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = local_node->v4_sock_meta; m; m = m->next) {
if (m->id == a->socket_id && (m->config_type == CFG_SERVER_TYPE_PUBLIC || m->config_type == CFG_SERVER_TYPE_UNKNOWN
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT)) {
out_addrs[direct_count].type = a->type; memcpy(out_addrs[direct_count].ip, a->addr, 4);
out_addrs[direct_count].port = htons(a->port); out_addrs[direct_count].socket_id = a->socket_id;
direct_count++; break;
}
}
}
if (direct_count == 0) return;
size_t pkt_size = CONN_MGR_DIRECT_REQ_HDR_SIZE + (size_t)direct_count * 8;
uint8_t* pkt = u_malloc(pkt_size);
if (!pkt) return;
struct CONN_MGR_DIRECT_REQ* req = (struct CONN_MGR_DIRECT_REQ*)pkt;
req->cmd = ETCP_RT_ID_CONN_MGR;
req->subcmd = CONN_MGR_SUBCMD_DIRECT_REQ;
req->request_id = 0;
req->addr_count = direct_count;
uint8_t* p = pkt + CONN_MGR_DIRECT_REQ_HDR_SIZE;
for (uint8_t i = 0; i < direct_count; i++) {
*p++ = out_addrs[i].type;
memcpy(p, out_addrs[i].ip, 4); p += 4;
memcpy(p, &out_addrs[i].port, 2); p += 2;
*p++ = out_addrs[i].socket_id;
}
struct ll_entry* qe = queue_entry_new(0);
if (!qe) { u_free(pkt); return; }
qe->dgram = pkt; qe->len = (uint16_t)pkt_size;
etcp_route_send(entry->conn->instance, entry->group_id, entry->node_id, qe, 1);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] REVERSE: sent DIRECT_REQ to 0x%016llx with %u addrs",
(unsigned long long)entry->node_id, direct_count);
}
/* ═══════════ REVERSE polling: проверка incoming conn ═══════════ */
static void ncd_reverse_poll_cb(void* arg) {
struct ncd_entry* entry = (struct ncd_entry*)arg;
if (!entry || entry->up || entry->timed_out || entry->conn->fin_wait) return;
struct ETCP_CONN* incoming = instance_find_conn(entry->conn->instance, entry->node_id);
if (incoming && incoming->state == 1 && incoming->peer_node_id == entry->node_id) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] REVERSE: incoming conn found node=0x%016llx", (unsigned long long)entry->node_id);
struct ETCP_CONN* old_conn = entry->conn;
entry->conn = incoming;
ncd_event_dispatch(entry, NCD_EVENT_UP);
if (old_conn && old_conn != incoming) etcp_connection_close(old_conn);
} else {
entry->reverse_poll_timer = uasync_set_timeout(entry->ua, 1000, entry, ncd_reverse_poll_cb, "ncd_rev_poll");
}
}
/* ═══════════ единая диспетчеризация событий ═══════════ */
static void ncd_event_dispatch(struct ncd_entry* entry, enum ncd_event event) {
@ -259,7 +156,6 @@ static void ncd_event_dispatch(struct ncd_entry* entry, enum ncd_event 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; }
if (entry->reverse_poll_timer) { uasync_cancel_timeout(entry->ua, entry->reverse_poll_timer); entry->reverse_poll_timer = NULL; }
entry->up = 1;
break;
case NCD_EVENT_DOWN:
@ -316,16 +212,6 @@ static void ncd_connect_timeout_cb(void* arg) {
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);
/* проверяем incoming conn (REVERSE) */
struct ETCP_CONN* incoming = instance_find_conn(entry->conn->instance, entry->node_id);
if (incoming && incoming->state == 1 && incoming->peer_node_id == entry->node_id) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] connect timeout but incoming found, UP");
struct ETCP_CONN* old_conn = entry->conn;
entry->conn = incoming;
ncd_event_dispatch(entry, NCD_EVENT_UP);
if (old_conn && old_conn != incoming) etcp_connection_close(old_conn);
return;
}
ncd_event_dispatch(entry, NCD_EVENT_TIMEOUT);
}
@ -415,7 +301,7 @@ static void ncd_init_control_binding(struct UTUN_INSTANCE* inst) {
/* ═══════════ API ═══════════ */
int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t group_id,
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) {
if (!inst || !out_handle) return NCD_ERR;
@ -448,7 +334,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t
/* 2. Ищем conn через instance_find_conn (входящее / созданное etcp_connect) */
struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (conn) {
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);
@ -456,7 +342,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t
}
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->group_id = group_id; entry->conn = conn; entry->ua = inst->ua;
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);
@ -524,7 +410,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t
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->group_id = group_id; entry->conn = conn; entry->ua = inst->ua; entry->up = 0;
entry->node_id = node_id; entry->conn = conn; entry->ua = inst->ua; entry->up = 0;
ncd_registry_add(entry);
struct NODE_CONN_DIRECT* h = u_calloc(1, sizeof(*h));
@ -541,18 +427,8 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t
int link_count = ncd_create_links(entry, ni);
if (link_count == 0 && group_id != 0) {
struct TOPO_GROUP* grp = topo_groups_find(inst->topo_groups, group_id);
struct TOPO_GROUP_NODE* local_nq = grp ? grp->local_node : NULL;
struct TOPO_NODE* local_node = local_nq ? topo_node_registry_find(inst->topo_groups, local_nq->node_id) : NULL;
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] no links, trying REVERSE node=0x%016llx group=0x%016llx",
(unsigned long long)node_id, (unsigned long long)group_id);
ncd_try_reverse(entry, local_node, ni);
entry->reverse_poll_timer = uasync_set_timeout(inst->ua, 1000, entry, ncd_reverse_poll_cb, "ncd_rev_poll");
}
if (link_count == 0) {
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] no links created for node=0x%016llx (will wait for incoming)", (unsigned long long)node_id);
}
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",
@ -561,6 +437,133 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t
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) {
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);
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(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; }
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 (временный, не владеем) */
{ struct ETCP_CONN* conn = etcp_connection_create(inst, NULL);
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(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; }
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);
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;
@ -578,7 +581,6 @@ void node_conn_direct_close(struct NODE_CONN_DIRECT* h) {
if (entry->handle_count <= 0) {
if (entry->connect_timer) { uasync_cancel_timeout(entry->ua, entry->connect_timer); entry->connect_timer = NULL; }
if (entry->reverse_poll_timer) { uasync_cancel_timeout(entry->ua, entry->reverse_poll_timer); entry->reverse_poll_timer = NULL; }
if (conn && conn->state != 2) {
/* устанавливаем fin_wait, отправляем CLOSE, ставим таймер */
conn->fin_wait = 1;
@ -607,63 +609,6 @@ void node_conn_direct_close(struct NODE_CONN_DIRECT* h) {
u_free(h);
}
int node_conn_direct_adopt(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn,
ncd_callback cb, void* cb_arg,
struct NODE_CONN_DIRECT** out_handle) {
if (!inst || !conn || !out_handle) return NCD_ERR;
*out_handle = NULL;
ncd_init_control_binding(inst);
uint64_t node_id = conn->peer_node_id;
/* снять fin_wait если был */
if (conn->fin_wait) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] adopt clearing fin_wait 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 = ncd_registry_find(node_id);
if (entry) {
/* entry уже есть — просто добавляем handle */
if (entry->fin_wait_timer) { uasync_cancel_timeout(entry->ua, entry->fin_wait_timer); entry->fin_wait_timer = NULL; }
if (entry->conn != conn) {
if (entry->conn && entry->conn->state != 2) etcp_connection_close(entry->conn);
entry->conn = conn;
}
struct NODE_CONN_DIRECT* h = u_calloc(1, sizeof(*h));
if (!h) 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 (conn->state == 1 && entry->up) uasync_call_soon(inst->ua, h, ncd_deliver_up_cb);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] adopt REUSED node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
return NCD_REUSED;
}
entry = u_calloc(1, sizeof(*entry));
if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] adopt 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);
struct NODE_CONN_DIRECT* h = u_calloc(1, sizeof(*h));
if (!h) { ncd_registry_remove(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] adopt NEW node=0x%016llx conn=%p handles=%d", (unsigned long long)node_id, (void*)conn, entry->handle_count);
return NCD_REUSED;
}
struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h) {
if (!h || !h->entry) return NULL;
return h->entry->conn;

16
src/transport_layer/node_conn_direct.h

@ -5,6 +5,7 @@
struct UTUN_INSTANCE;
struct ETCP_CONN;
struct TOPO_NODE;
/* ─── Return codes ─── */
#define NCD_NEW 0
@ -26,24 +27,19 @@ typedef void (*ncd_callback)(struct NODE_CONN_DIRECT* h, enum ncd_event event, v
/*
* Открыть handle к узлу по node_id.
* Node info — из node_registry, fallback SQLite.
* group_id — для REVERSE-подключения через etcp_route_send (0 = без REVERSE).
* cb — вызывается при каждом изменении статуса / таймауте.
* Если conn уже готов — cb(h, NCD_EVENT_UP) через uasync_call_soon (не синхронно).
*
* Возвращает NCD_NEW / NCD_REUSED / NCD_ERR.
*/
int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id, uint64_t group_id,
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);
/*
* Обернуть уже существующий ETCP_CONN в handle (crypto+links уже готовы).
* Не создаёт линки, не делает crypto, не ставит connect_timer.
* Возвращает NCD_REUSED / NCD_ERR.
*/
int node_conn_direct_adopt(struct UTUN_INSTANCE* inst, struct ETCP_CONN* conn,
ncd_callback cb, void* cb_arg,
struct NODE_CONN_DIRECT** out_handle);
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);
void node_conn_direct_close(struct NODE_CONN_DIRECT* h);

135
tests/test_node_conn_direct.c

@ -7,10 +7,17 @@
* 2 — poll → init_cb → NCD_EVENT_UP на обоих
* 3 — open на ready conn → NCD_REUSED + async cb
* 4 — ещё open → NCD_REUSED
* 5 — close h3,h2 → conn жив
* 6 — close h1,h0 (last) → conn закрыт (fin_wait → CLOSE → DOWN)
* 7 — NCD_ERR unknown node
* 8 — inject unreachable → NCD_NEW → TIMEOUT
* 5 — close h3,h2 → conn жив (h0,h1 остаются)
* 6a — bidirectional: inject A→B, открыть B→A
* 6b — poll: A и B UP. A закрывает новый handle → CLOSE к B
* 6c — B имеет handle → KEEP_ALIVE. A переоткрывает → NCD_REUSED
* 6d — poll: A get_conn OK после KEEP_ALIVE. Cleanup B-side.
* 6e — open во время fin_wait: A закрывает h1, сразу же переоткрывает → NCD_REUSED
* 6f — poll: UP на переоткрытом handle. Cleanup.
* 6g — CLOSE от peer когда нет handle'ов: A закрывает → CLOSE к B, B без handles → conn dead → entry удалён → NCD_NEW при переоткрытии
* 6 — close h1,h0 (last) → conn закрыт (fin_wait → CLOSE → reinit)
* 7 — NCD_ERR unknown node
* 8 — inject unreachable → NCD_NEW → TIMEOUT
*/
#include <stdio.h>
#include <stdlib.h>
@ -50,11 +57,13 @@ static void *g_ttimer = NULL;
static char tdir[] = "/tmp/utun_ncd_XXXXXX";
static char ca[256], cb[256];
static int pa = 0, pb = 0;
static uint64_t nid_b = 0;
static uint8_t pubkey_b[SC_PUBKEY_SIZE];
static uint64_t nid_a = 0, nid_b = 0;
static uint8_t pubkey_a[SC_PUBKEY_SIZE], pubkey_b[SC_PUBKEY_SIZE];
static struct NODE_CONN_DIRECT *gh[4];
static volatile int g_up[4], g_down[4], g_tout[4];
static struct NODE_CONN_DIRECT *gh[8];
static volatile int g_up[8], g_down[8], g_tout[8];
static struct NODE_CONN_DIRECT *gh_b[2];
static volatile int g_up_b[2], g_down_b[2];
static void ncd_cb(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) {
int idx = (int)(intptr_t)arg;
@ -64,14 +73,26 @@ static void ncd_cb(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg)
fprintf(stderr, " cb[%d]: event=%d (up=%d down=%d tout=%d)\n", idx, (int)event, g_up[idx], g_down[idx], g_tout[idx]); fflush(stderr);
}
static void ncd_cb_b(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) {
int idx = (int)(intptr_t)arg;
if (event == NCD_EVENT_UP) g_up_b[idx]++;
else if (event == NCD_EVENT_DOWN) g_down_b[idx]++;
fprintf(stderr, " cb_b[%d]: event=%d (up=%d down=%d)\n", idx, (int)event, g_up_b[idx], g_down_b[idx]); fflush(stderr);
}
static void to_cb(void* arg) { (void)arg; fprintf(stderr, "GLOBAL TIMEOUT phase=%d\n", g_phase); g_result = 2; }
static int wf(const char* p, const char* f, ...) { va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; }
static void fail(const char* msg) { fprintf(stderr, "FAIL[%d]: %s\n", g_phase, msg); fflush(stderr); g_result = 2; }
static int inject_node(struct UTUN_INSTANCE* inst, uint64_t nid, const uint8_t pk[SC_PUBKEY_SIZE], uint16_t port);
static void t1(void* arg); static void t2(void* arg); static void t3(void* arg);
static void t4(void* arg); static void t5(void* arg); static void t6(void* arg);
static void t7(void* arg); static void t8(void* arg); static void t9(void* arg);
static void t4(void* arg); static void t5(void* arg);
static void t6a(void* arg); static void t6b(void* arg); static void t6c(void* arg);
static void t6d(void* arg); static void t6e(void* arg); static void t6f(void* arg);
static void t6g(void* arg);
static void t6(void* arg); static void t7(void* arg);
static void t8(void* arg); static void t9(void* arg);
static void t_done(void* arg);
static void t1(void* arg) {
@ -79,11 +100,11 @@ static void t1(void* arg) {
fprintf(stderr, "\n=== P1: two opens before poll ===\n"); fflush(stderr);
g_a->etcp_connect_timeout_tb = TIMEOUT_TB;
int r = node_conn_direct_open(g_a, nid_b, 0, ncd_cb, (void*)0, &gh[0]);
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)0, &gh[0]);
fprintf(stderr, " open#0 → %s\n", r == NCD_NEW ? "NCD_NEW" : r == NCD_REUSED ? "NCD_REUSED" : "ERR");
if (r != NCD_NEW) { fail("expected NCD_NEW"); return; }
r = node_conn_direct_open(g_a, nid_b, 0, ncd_cb, (void*)1, &gh[1]);
r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)1, &gh[1]);
fprintf(stderr, " open#1 → %s\n", r == NCD_NEW ? "NCD_NEW" : r == NCD_REUSED ? "NCD_REUSED" : "ERR");
if (r != NCD_REUSED) { fail("expected NCD_REUSED"); return; }
@ -104,7 +125,7 @@ static void t2(void* arg) {
static void t3(void* arg) {
(void)arg; g_phase = 3;
fprintf(stderr, "\n=== P3: open on ready conn ===\n"); fflush(stderr);
int r = node_conn_direct_open(g_a, nid_b, 0, ncd_cb, (void*)2, &gh[2]);
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)2, &gh[2]);
fprintf(stderr, " open#2 → %s\n", r == NCD_REUSED ? "NCD_REUSED" : "?");
if (r != NCD_REUSED) { fail("expected NCD_REUSED"); return; }
uasync_call_soon(ua, NULL, t4);
@ -114,7 +135,7 @@ static void t4(void* arg) {
(void)arg; g_phase = 4;
if (g_up[2] == 0) { uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t4, "t4"); return; }
fprintf(stderr, "\n=== P4: async cb OK, open h3 ===\n"); fflush(stderr);
int r = node_conn_direct_open(g_a, nid_b, 0, ncd_cb, (void*)3, &gh[3]);
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)3, &gh[3]);
fprintf(stderr, " open#3 → %s\n", r == NCD_REUSED ? "NCD_REUSED" : "?");
if (r != NCD_REUSED) { fail("expected NCD_REUSED"); return; }
uasync_call_soon(ua, NULL, t5);
@ -129,7 +150,71 @@ static void t5(void* arg) {
node_conn_direct_close(gh[2]); gh[2] = NULL;
if (node_conn_direct_get_conn(gh[1]) == NULL) { fail("conn died after h2 close"); return; }
fprintf(stderr, " OK\n"); fflush(stderr);
uasync_call_soon(ua, NULL, t6);
uasync_call_soon(ua, NULL, t6a);
}
/* ═══════════ Scenario 1: KEEP_ALIVE handshake (conn ещё жив после P1-P5) ═══════════ */
static void t6a(void* arg) {
(void)arg; g_phase = 16;
fprintf(stderr, "\n=== P6a: bidirectional — B opens handle to A ===\n"); fflush(stderr);
int r = node_conn_direct_open(g_b, nid_a, ncd_cb_b, (void*)0, &gh_b[0]);
if (r != NCD_REUSED) { fail("P6a: B open expected NCD_REUSED (incoming conn exists)"); return; }
fprintf(stderr, " B open → NCD_REUSED\n");
r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)4, &gh[4]);
if (r != NCD_REUSED) { fail("P6a: A open(4) expected NCD_REUSED (conn alive)"); return; }
fprintf(stderr, " A open(4) → NCD_REUSED\n"); fflush(stderr);
uasync_call_soon(ua, NULL, t6b);
}
static void t6b(void* arg) {
(void)arg; g_phase = 17;
if (g_up[4] > 0 && g_up_b[0] > 0) {
fprintf(stderr, "\n=== P6b: both UP, A closes handle 4 → CLOSE to B ===\n"); fflush(stderr);
node_conn_direct_close(gh[4]); gh[4] = NULL;
uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)t6c, "t6b_d");
return;
}
uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t6b, "t6b");
}
static void t6c(void* arg) {
(void)arg; g_phase = 18;
fprintf(stderr, "\n=== P6c: KEEP_ALIVE round-trip, A re-opens ===\n"); fflush(stderr);
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)5, &gh[5]);
if (r != NCD_REUSED) { fail("P6c: expected NCD_REUSED (KEEP_ALIVE saved conn)"); return; }
fprintf(stderr, " A re-open(5) → NCD_REUSED\n"); fflush(stderr);
uasync_call_soon(ua, NULL, t6d);
}
static void t6d(void* arg) {
(void)arg; g_phase = 19;
if (g_up[5] == 0) { uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t6d, "t6d"); return; }
if (node_conn_direct_get_conn(gh[5]) == NULL) { fail("P6d: get_conn NULL after KEEP_ALIVE"); return; }
fprintf(stderr, "\n=== P6d: KEEP_ALIVE works, get_conn OK ===\n"); fflush(stderr);
node_conn_direct_close(gh[5]); gh[5] = NULL;
node_conn_direct_close(gh_b[0]); gh_b[0] = NULL;
uasync_set_timeout(ua, 300, NULL, (timeout_callback_t)t6e, "t6d_d");
}
/* ═══════════ Scenario 2: Open during fin_wait ═══════════ */
static void t6e(void* arg) {
(void)arg; g_phase = 20;
fprintf(stderr, "\n=== P6e: close h1 → fin_wait, immediately re-open ===\n"); fflush(stderr);
node_conn_direct_close(gh[1]); gh[1] = NULL;
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)4, &gh[4]);
if (r != NCD_REUSED) { fail("P6e: expected NCD_REUSED (entry found, fin_wait cleared)"); return; }
fprintf(stderr, " re-open → NCD_REUSED\n"); fflush(stderr);
uasync_call_soon(ua, NULL, t6f);
}
static void t6f(void* arg) {
(void)arg; g_phase = 21;
if (g_up[4] == 0) { uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t6f, "t6f"); return; }
fprintf(stderr, "\n=== P6f: re-opened handle got UP, Scenario 2 OK ===\n"); fflush(stderr);
node_conn_direct_close(gh[4]); gh[4] = NULL;
uasync_set_timeout(ua, 300, NULL, (timeout_callback_t)t6, "t6f_d");
}
static void t6(void* arg) {
@ -139,14 +224,14 @@ static void t6(void* arg) {
if (node_conn_direct_get_conn(gh[0]) == NULL) { fail("conn died before last handle"); return; }
node_conn_direct_close(gh[0]); gh[0] = NULL;
fprintf(stderr, " OK\n"); fflush(stderr);
uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)t7, "t6d");
uasync_set_timeout(ua, 3000, NULL, (timeout_callback_t)t7, "t6d");
}
static void t7(void* arg) {
(void)arg; g_phase = 7;
fprintf(stderr, "\n=== P7: NCD_ERR unknown node ===\n"); fflush(stderr);
struct NODE_CONN_DIRECT* hx = NULL;
int r = node_conn_direct_open(g_a, 0xDEADBEEF00000001ULL, 0, ncd_cb, (void*)99, &hx);
int r = node_conn_direct_open(g_a, 0xDEADBEEF00000001ULL, ncd_cb, (void*)99, &hx);
fprintf(stderr, " open(unknown) → %s\n", r == NCD_ERR ? "NCD_ERR" : "?");
if (r != NCD_ERR) { fail("expected NCD_ERR"); return; }
if (hx != NULL) { fail("handle not NULL"); return; }
@ -176,7 +261,7 @@ static void t8(void* arg) {
memset((void*)g_tout, 0, sizeof(g_tout));
g_a->etcp_connect_timeout_tb = SHORT_TO_TB;
int r = node_conn_direct_open(g_a, fake_id, 0, ncd_cb, (void*)0, &gh[0]);
int r = node_conn_direct_open(g_a, fake_id, ncd_cb, (void*)0, &gh[0]);
fprintf(stderr, " open(unreachable) → %s\n", r == NCD_NEW ? "NCD_NEW" : "?");
if (r != NCD_NEW) { fail("expected NCD_NEW"); return; }
@ -216,8 +301,8 @@ static int setup(void) {
test_mkdtemp(tdir);
int base = 49000 + (getpid() % 10000); pa = base; pb = base + 1;
snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir);
wf(ca, "[global]\ntun_ip=10.97.0.1/24\ntun_ifname=tun95\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pa);
wf(cb, "[global]\ntun_ip=10.97.0.2/24\ntun_ifname=tun94\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pb);
wf(ca, "[global]\ntun_ip=10.97.0.1/24\ntun_ifname=tun95\nkeepalive_timeout=100\nkeepalive_interval=10\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pa);
wf(cb, "[global]\ntun_ip=10.97.0.2/24\ntun_ifname=tun94\nkeepalive_timeout=100\nkeepalive_interval=10\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pb);
config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb);
return 0;
}
@ -231,17 +316,21 @@ int main(void) {
g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb);
if (!g_a || !g_b) { g_result = 2; goto done; }
nid_b = g_b->node_id; memcpy(pubkey_b, g_b->my_keys.public_key, SC_PUBKEY_SIZE);
fprintf(stderr, "B: 0x%016llx port=%d\n", (unsigned long long)nid_b, pb); fflush(stderr);
nid_a = g_a->node_id; memcpy(pubkey_a, g_a->my_keys.public_key, SC_PUBKEY_SIZE);
fprintf(stderr, "A: 0x%016llx port=%d B: 0x%016llx port=%d\n",
(unsigned long long)nid_a, pa, (unsigned long long)nid_b, pb); fflush(stderr);
utun_instance_init(g_a); utun_instance_init(g_b);
if (inject_node(g_a, nid_b, pubkey_b, pb) != 0) { fail("inject_node"); goto done; }
fprintf(stderr, "injected B into A registry\n"); fflush(stderr);
if (inject_node(g_b, nid_a, pubkey_a, pa) != 0) { fail("inject_node A→B"); goto done; }
fprintf(stderr, "injected B→A and A→B registries\n"); fflush(stderr);
uasync_call_soon(ua, NULL, t1);
g_ttimer = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to");
{ uint64_t st = get_time_tb(); while (!g_result && (int)(get_time_tb() - st) < TIMEOUT_TB + 50000) uasync_poll(ua, POLL_MS); }
if (g_result == 0) g_result = 2;
done:
if (g_ttimer && ua) uasync_cancel_timeout(ua, g_ttimer);
for (int i = 0; i < 4; i++) if (gh[i]) node_conn_direct_close(gh[i]);
for (int i = 0; i < 8; i++) if (gh[i]) node_conn_direct_close(gh[i]);
for (int i = 0; i < 2; i++) if (gh_b[i]) node_conn_direct_close(gh_b[i]);
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); }
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); }
if (ua) uasync_destroy(ua, 0);

Loading…
Cancel
Save