diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index bca1183e..32279f46 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/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; diff --git a/src/transport_layer/node_conn_direct.h b/src/transport_layer/node_conn_direct.h index 90129503..d65ec03d 100644 --- a/src/transport_layer/node_conn_direct.h +++ b/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); diff --git a/tests/test_node_conn_direct.c b/tests/test_node_conn_direct.c index 232629ef..87ec5ca7 100644 --- a/tests/test_node_conn_direct.c +++ b/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 #include @@ -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);