Browse Source
- test_etcp_router_reconnect, test_etcp_connect, test_stcp_traffic: replace u_calloc with memory_pool_alloc for TOPO_SOCKMETA4/TOPO_ADDR4 - db_sync.c: fix req[14] -> req[15] — SEND_DATA request header is 15 bytes (type:1 + from:4 + count:2 + vp:4 + want_from:4)topo_upd
52 changed files with 1796 additions and 585 deletions
@ -0,0 +1,670 @@
|
||||
/*
|
||||
* node_conn_direct.c — handle-based ETCP connection layer |
||||
* |
||||
* Тонкая прослойка над 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" |
||||
#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 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; |
||||
|
||||
struct NODE_CONN_DIRECT* handles; /* связный список всех handle'ов */ |
||||
|
||||
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 уже сработал */ |
||||
|
||||
struct ncd_entry* next; /* цепочка статического реестра */ |
||||
}; |
||||
|
||||
struct NODE_CONN_DIRECT { |
||||
struct ncd_entry* entry; |
||||
ncd_callback cb; |
||||
void* cb_arg; |
||||
struct NODE_CONN_DIRECT* next; |
||||
}; |
||||
|
||||
static struct ncd_entry* g_ncd_registry; |
||||
static uint8_t g_ncd_control_bound; /* 1 = etcp_bind(ETCP_RT_ID_NCD_CONTROL) уже сделан */ |
||||
|
||||
/* ─── forward declarations ─── */ |
||||
static void ncd_event_dispatch(struct ncd_entry* entry, enum ncd_event event); |
||||
static void ncd_connect_timeout_cb(void* arg); |
||||
static void ncd_deliver_up_cb(void* arg); |
||||
|
||||
/* ═══════════ реестр ═══════════ */ |
||||
|
||||
static struct ncd_entry* ncd_registry_find(uint64_t node_id) { |
||||
struct ncd_entry* e = g_ncd_registry; |
||||
while (e) { if (e->node_id == node_id) return e; e = e->next; } |
||||
return NULL; |
||||
} |
||||
static void ncd_registry_add(struct ncd_entry* entry) { |
||||
entry->next = g_ncd_registry; g_ncd_registry = entry; |
||||
} |
||||
static void ncd_registry_remove(struct ncd_entry* entry) { |
||||
struct ncd_entry** pp = &g_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_CONN* conn = entry->conn; |
||||
struct UTUN_INSTANCE* inst = conn->instance; |
||||
int link_count = 0; |
||||
|
||||
/* ---- IPv4 ---- */ |
||||
{ |
||||
struct ETCP_SOCKET* socks[64]; int sock_count = 0; |
||||
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) |
||||
if (s->local_addr.ss_family == AF_INET && s->type != CFG_SERVER_TYPE_PRIVATE) |
||||
socks[sock_count++] = s; |
||||
|
||||
if (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)); |
||||
if (etcp_link_new(conn, socks[rr++ % sock_count], &sa, 0)) link_count++; |
||||
} |
||||
} |
||||
} |
||||
|
||||
/* ---- IPv6 ---- */ |
||||
{ |
||||
struct ETCP_SOCKET* socks[64]; int sock_count = 0; |
||||
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) |
||||
if (s->local_addr.ss_family == AF_INET6 && s->type != CFG_SERVER_TYPE_PRIVATE) |
||||
socks[sock_count++] = s; |
||||
|
||||
if (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 sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); |
||||
if (etcp_link_new(conn, socks[rr++ % sock_count], &sa, 0)) link_count++; |
||||
} |
||||
} |
||||
} |
||||
|
||||
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) { |
||||
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; } |
||||
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: |
||||
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=%d node=0x%016llx handles=%d", |
||||
(int)event, (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_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; } |
||||
etcp_conn_remove_cbk(conn, ncd_init_cb, entry);; |
||||
etcp_conn_remove_cbk(conn, ncd_up_cb, entry);; |
||||
etcp_conn_remove_cbk(conn, ncd_down_cb, entry);; |
||||
etcp_connection_close(conn); |
||||
ncd_registry_remove(entry); |
||||
u_free(entry); |
||||
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); |
||||
/* проверяем 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); |
||||
} |
||||
|
||||
/* ═══════════ асинхронная доставка 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); |
||||
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(sender_id); |
||||
|
||||
switch (msg->subcmd) { |
||||
case NCD_SUBCMD_CLOSE: |
||||
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] recv CLOSE from 0x%016llx", (unsigned long long)sender_id); |
||||
if (!entry) { |
||||
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] no entry for CLOSE sender 0x%016llx, closing conn", (unsigned long long)sender_id); |
||||
if (conn) etcp_connection_close(conn); |
||||
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] closing conn for CLOSE from 0x%016llx", (unsigned long long)sender_id); |
||||
if (conn) etcp_connection_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; } |
||||
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 (g_ncd_control_bound) return; |
||||
etcp_bind(inst, ETCP_RT_ID_NCD_CONTROL, ncd_recv_control_handler); |
||||
g_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, uint64_t group_id, |
||||
ncd_callback cb, void* cb_arg, |
||||
struct NODE_CONN_DIRECT** out_handle) { |
||||
if (!inst || !out_handle) return NCD_ERR; |
||||
*out_handle = NULL; |
||||
ncd_init_control_binding(inst); |
||||
|
||||
/* 1. Ищем в реестре — conn уже есть */ |
||||
struct ncd_entry* entry = ncd_registry_find(node_id); |
||||
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) { |
||||
/* снять 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->group_id = group_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] 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 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; |
||||
} |
||||
|
||||
conn = etcp_connection_create(inst, NULL); |
||||
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->group_id = group_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] alloc handle failed node=0x%016llx", (unsigned long long)node_id); |
||||
ncd_registry_remove(entry); u_free(entry); etcp_connection_close(conn); topo_node_registry_unref(inst->topo_groups, ni->node_id); return NCD_ERR; } |
||||
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 && 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); |
||||
} |
||||
|
||||
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; |
||||
} |
||||
|
||||
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 (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; |
||||
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 { |
||||
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);; |
||||
if (conn) etcp_connection_close(conn); |
||||
ncd_registry_remove(entry); |
||||
u_free(entry); |
||||
} |
||||
} |
||||
|
||||
h->entry = NULL; |
||||
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; |
||||
} |
||||
@ -0,0 +1,67 @@
|
||||
#ifndef NODE_CONN_DIRECT_H |
||||
#define NODE_CONN_DIRECT_H |
||||
|
||||
#include <stdint.h> |
||||
|
||||
struct UTUN_INSTANCE; |
||||
struct ETCP_CONN; |
||||
|
||||
/* ─── Return codes ─── */ |
||||
#define NCD_NEW 0 |
||||
#define NCD_REUSED 1 |
||||
#define NCD_ERR -1 |
||||
|
||||
/* ─── Opaque handle ─── */ |
||||
struct NODE_CONN_DIRECT; |
||||
|
||||
/* ─── Callback events (единый callback на handle) ─── */ |
||||
enum ncd_event { |
||||
NCD_EVENT_UP = 0, |
||||
NCD_EVENT_DOWN = 1, |
||||
NCD_EVENT_TIMEOUT = 2, |
||||
}; |
||||
|
||||
typedef void (*ncd_callback)(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg); |
||||
|
||||
/*
|
||||
* Открыть 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, |
||||
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); |
||||
|
||||
void node_conn_direct_close(struct NODE_CONN_DIRECT* h); |
||||
|
||||
struct ETCP_CONN* node_conn_direct_get_conn(struct NODE_CONN_DIRECT* h); |
||||
|
||||
/* ─── Протокол CLOSE / KEEP_ALIVE (ETCP_RT_ID_NCD_CONTROL = 0x12) ─── */ |
||||
|
||||
#define NCD_SUBCMD_CLOSE 0x01 |
||||
#define NCD_SUBCMD_KEEP_ALIVE 0x02 |
||||
|
||||
#pragma pack(push, 1) |
||||
struct ncd_control_msg { |
||||
uint8_t cmd; /* ETCP_RT_ID_NCD_CONTROL */ |
||||
uint8_t subcmd; /* NCD_SUBCMD_CLOSE / KEEP_ALIVE */ |
||||
uint64_t node_id; /* node_id отправителя */ |
||||
}; |
||||
#pragma pack(pop) |
||||
|
||||
#define NCD_CONTROL_MSG_SIZE sizeof(struct ncd_control_msg) |
||||
|
||||
#endif |
||||
@ -0,0 +1,251 @@
|
||||
/**
|
||||
* @file test_node_conn_direct.c |
||||
* @brief Тест node_conn_direct с 2 инстансами. |
||||
* |
||||
* Фазы: |
||||
* 1 — два open ДО eventloop → NCD_NEW + NCD_REUSED(pending) |
||||
* 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 |
||||
*/ |
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include <stdarg.h> |
||||
#include "../lib/platform_compat.h" |
||||
#include "test_utils.h" |
||||
#ifndef _WIN32 |
||||
#include <unistd.h> |
||||
#endif |
||||
|
||||
#include "etcp.h" |
||||
#include "etcp_connections.h" |
||||
#include "node_conn_direct.h" |
||||
#include "../src/config_parser.h" |
||||
#include "../src/config_updater.h" |
||||
#include "../src/utun_instance.h" |
||||
#include "topo_group.h" |
||||
#include "topo_node.h" |
||||
#include "secure_channel.h" |
||||
#include "../lib/u_async.h" |
||||
#include "../lib/ll_queue.h" |
||||
#include "../lib/memory_pool.h" |
||||
#include "../lib/debug_config.h" |
||||
#include "../lib/mem.h" |
||||
|
||||
#define TIMEOUT_TB 600000 |
||||
#define SHORT_TO_TB 50000 |
||||
#define POLL_MS 5 |
||||
|
||||
static struct UTUN_INSTANCE *g_a = NULL, *g_b = NULL; |
||||
static struct UASYNC *ua = NULL; |
||||
static volatile int g_result = 0; |
||||
static int g_phase = 0; |
||||
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 struct NODE_CONN_DIRECT *gh[4]; |
||||
static volatile int g_up[4], g_down[4], g_tout[4]; |
||||
|
||||
static void ncd_cb(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { |
||||
int idx = (int)(intptr_t)arg; |
||||
if (event == NCD_EVENT_UP) g_up[idx]++; |
||||
else if (event == NCD_EVENT_DOWN) g_down[idx]++; |
||||
else if (event == NCD_EVENT_TIMEOUT) g_tout[idx]++; |
||||
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 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 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 t_done(void* arg); |
||||
|
||||
static void t1(void* arg) { |
||||
(void)arg; g_phase = 1; |
||||
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]); |
||||
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]); |
||||
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; } |
||||
|
||||
uasync_call_soon(ua, NULL, t2); |
||||
} |
||||
|
||||
static void t2(void* arg) { |
||||
(void)arg; g_phase = 2; |
||||
if (g_up[0] > 0 && g_up[1] > 0) { |
||||
if (node_conn_direct_get_conn(gh[0]) == NULL) { fail("get_conn h0 NULL"); return; } |
||||
if (node_conn_direct_get_conn(gh[1]) == NULL) { fail("get_conn h1 NULL"); return; } |
||||
fprintf(stderr, "\n=== P2: both up, get_conn OK ===\n"); fflush(stderr); |
||||
uasync_call_soon(ua, NULL, t3); return; |
||||
} |
||||
uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t2, "t2"); |
||||
} |
||||
|
||||
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]); |
||||
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); |
||||
} |
||||
|
||||
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]); |
||||
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); |
||||
} |
||||
|
||||
static void t5(void* arg) { |
||||
(void)arg; g_phase = 5; |
||||
if (g_up[3] == 0) { uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t5, "t5"); return; } |
||||
fprintf(stderr, "\n=== P5: close h3, h2 → conn stays ===\n"); fflush(stderr); |
||||
node_conn_direct_close(gh[3]); gh[3] = NULL; |
||||
if (node_conn_direct_get_conn(gh[2]) == NULL) { fail("conn died after h3 close"); return; } |
||||
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); |
||||
} |
||||
|
||||
static void t6(void* arg) { |
||||
(void)arg; g_phase = 6; |
||||
fprintf(stderr, "\n=== P6: close last → conn closed ===\n"); fflush(stderr); |
||||
node_conn_direct_close(gh[1]); gh[1] = NULL; |
||||
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"); |
||||
} |
||||
|
||||
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); |
||||
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; } |
||||
fprintf(stderr, " OK\n"); fflush(stderr); |
||||
uasync_call_soon(ua, NULL, t8); |
||||
} |
||||
|
||||
static void t8(void* arg) { |
||||
(void)arg; g_phase = 8; |
||||
fprintf(stderr, "\n=== P8: unreachable node → timeout ===\n"); fflush(stderr); |
||||
|
||||
struct TOPO_NODE* ni = u_calloc(1, sizeof(struct TOPO_NODE)); |
||||
if (!ni) { fail("alloc ni"); return; } |
||||
uint64_t fake_id = 0xF000000000000001ULL; |
||||
ni->group_ref_count = 0; ni->node_id = fake_id; ni->ver = 0; |
||||
memcpy(ni->public_key, pubkey_b, SC_PUBKEY_SIZE); |
||||
|
||||
struct TOPO_SOCKMETA4* sm = memory_pool_alloc(g_a->topo_groups->v4_sock_meta_pool); |
||||
if (sm) { sm->id = 0; sm->config_type = CFG_SERVER_TYPE_PUBLIC; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; } |
||||
|
||||
struct TOPO_ADDR4* addr = memory_pool_alloc(g_a->topo_groups->v4_addr_pool); |
||||
if (addr) { addr->addr[0]=127; addr->addr[1]=0; addr->addr[2]=0; addr->addr[3]=1; addr->port=1; addr->type=TOPO_ADDR_INTERFACE; addr->socket_id=0; addr->protocol=TOPO_PROTO_UDP; addr->next=ni->v4_addrs; ni->v4_addrs=addr; } |
||||
|
||||
if (!topo_node_registry_store(g_a->topo_groups, ni)) { fail("registry acquire"); return; } |
||||
|
||||
memset((void*)g_up, 0, sizeof(g_up)); |
||||
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]); |
||||
fprintf(stderr, " open(unreachable) → %s\n", r == NCD_NEW ? "NCD_NEW" : "?"); |
||||
if (r != NCD_NEW) { fail("expected NCD_NEW"); return; } |
||||
|
||||
uasync_call_soon(ua, NULL, t9); |
||||
} |
||||
|
||||
static void t9(void* arg) { |
||||
(void)arg; g_phase = 9; |
||||
if (g_tout[0] == 0) { uasync_set_timeout(ua, POLL_MS, NULL, (timeout_callback_t)t9, "t9"); return; } |
||||
if (g_up[0] > 0) { fail("up fired, expected only timeout"); return; } |
||||
fprintf(stderr, "\n=== P9: timeout_cb fired, no up ===\n"); fflush(stderr); |
||||
node_conn_direct_close(gh[0]); gh[0] = NULL; |
||||
uasync_call_soon(ua, NULL, t_done); |
||||
} |
||||
|
||||
static void t_done(void* arg) { |
||||
(void)arg; g_phase = 99; |
||||
fprintf(stderr, "\n=== ALL PASSED ===\n"); fflush(stderr); |
||||
g_result = 1; |
||||
} |
||||
|
||||
static int inject_node(struct UTUN_INSTANCE* inst, uint64_t nid, const uint8_t pk[SC_PUBKEY_SIZE], uint16_t port) { |
||||
if (!inst->topo_groups) return -1; |
||||
struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); |
||||
if (!ni) return -1; |
||||
ni->group_ref_count = 0; ni->node_id = nid; ni->ver = 0; |
||||
memcpy(ni->public_key, pk, SC_PUBKEY_SIZE); |
||||
struct TOPO_SOCKMETA4* sm = memory_pool_alloc(inst->topo_groups->v4_sock_meta_pool); |
||||
if (sm) { sm->id = 0; sm->config_type = CFG_SERVER_TYPE_PUBLIC; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; } |
||||
struct TOPO_ADDR4* addr = memory_pool_alloc(inst->topo_groups->v4_addr_pool); |
||||
if (addr) { addr->addr[0]=127; addr->addr[1]=0; addr->addr[2]=0; addr->addr[3]=1; addr->port=port; addr->type=TOPO_ADDR_INTERFACE; addr->socket_id=0; addr->protocol=TOPO_PROTO_UDP; addr->next=ni->v4_addrs; ni->v4_addrs=addr; } |
||||
if (!topo_node_registry_store(inst->topo_groups, ni)) { u_free(ni); return -1; } |
||||
return 0; |
||||
} |
||||
|
||||
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); |
||||
config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); |
||||
return 0; |
||||
} |
||||
static void cleanup(void) { test_unlink(ca); test_unlink(cb); test_rmdir(tdir); } |
||||
|
||||
int main(void) { |
||||
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); |
||||
utun_instance_set_tun_init_enabled(0); |
||||
if (setup() != 0) return 1; |
||||
ua = uasync_create(); |
||||
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); |
||||
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); |
||||
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]); |
||||
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); |
||||
cleanup(); |
||||
fprintf(stderr, "\n%s\n", g_result == 1 ? "=== TEST PASSED ===" : "=== TEST FAILED ==="); fflush(stderr); |
||||
return (g_result == 1) ? 0 : 1; |
||||
} |
||||
Loading…
Reference in new issue