From a7b20f268fb1697ee2e400902dbe9fcf03fa07de Mon Sep 17 00:00:00 2001 From: Evgeny Date: Mon, 29 Jun 2026 20:01:40 +0300 Subject: [PATCH] feat: conn_mgr, dummynet block-addr, indirect routing, alien nodes, nodeinfo refactor --- src/Makefile.am | 1 + src/conn_mgr.c | 1000 ++++++++++++++++++++++++++++++++++++ src/conn_mgr.h | 190 +++++++ src/dummynet.c | 35 ++ src/dummynet.h | 2 + src/etcp.h | 2 +- src/etcp_api.h | 3 + src/etcp_router.c | 15 +- src/msg_transport.c | 2 +- src/route_bgp.c | 25 +- src/route_bgp.h | 7 - src/route_node.c | 9 +- src/route_node.h | 13 + src/route_ping.c | 7 +- src/utun_instance.c | 6 + src/utun_instance.h | 2 + tests/Makefile.am | 6 + tests/test_conn_mgr.c | 136 +++++ tests/test_ipv6_sockets.c | 2 +- tests/test_nat_detection.c | 16 +- tests/test_nat_transport.c | 4 +- tests/test_route_ping.c | 6 +- 22 files changed, 1449 insertions(+), 40 deletions(-) create mode 100644 src/conn_mgr.c create mode 100644 src/conn_mgr.h create mode 100644 tests/test_conn_mgr.c diff --git a/src/Makefile.am b/src/Makefile.am index 69e38f7c..3d286e75 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -13,6 +13,7 @@ utun_CORE_SOURCES = \ route_node.c \ route_node_lmdb.c \ route_connectivity.c \ + conn_mgr.c \ routing.c \ tun_if.c \ tun_route.c \ diff --git a/src/conn_mgr.c b/src/conn_mgr.c new file mode 100644 index 00000000..b709af88 --- /dev/null +++ b/src/conn_mgr.c @@ -0,0 +1,1000 @@ +#include +#include +#ifdef _WIN32 +#include +#include +#else +#include +#endif +#include "../lib/platform_compat.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "utun_instance.h" +#include "etcp.h" +#include "etcp_connections.h" +#include "etcp_router.h" +#include "config_parser.h" +#include "secure_channel.h" +#include "route_node.h" +#include "route_bgp.h" +#include "route_ping.h" +#include "route_connectivity.h" +#include "conn_mgr.h" + +static struct CONN_MGR_ENTRY* cm_find_entry(struct CONN_MGR* mgr, uint64_t node_id); +static void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry); +static void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry); +static int cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry); +static void cm_start_local_scan(struct CONN_MGR_ENTRY* entry); +static void cm_idle_timer_cb(void* arg); +static void cm_deliver_result(struct CONN_MGR_ENTRY* entry, int result); +static void cm_bg_ping_timer_cb(void* arg); +static int cm_nat_compatible(struct ETCP_SOCKET* our, uint8_t target_meta_type); +static uint16_t cm_get_node_min_rtt(struct NODEINFO_Q* nq); +static int cm_has_direct_ip(struct NODEINFO_Q* nq); +static int cm_has_local_addr(struct NODEINFO_Q* nq); +static void cm_entry_destroy(struct CONN_MGR_ENTRY* entry); +static uint64_t cm_get_node_max_probe_time(struct NODEINFO_Q* nq); +static int cm_is_rtt_fresh(struct NODEINFO_Q* nq, uint64_t now_tb); + +static void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry); +static void cm_handle_direct_req(struct ETCP_CONN* conn, const uint8_t* data, size_t len); +static void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CONN_MGR_INTERM_EXCHANGE_REQ* req); +static void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len); +static void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len); +static void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id); +static void cm_direct_ready_cb(struct ETCP_CONN* conn, void* arg); +static void cm_reverse_ready_cb(struct ETCP_CONN* conn, void* arg); +static void cm_reverse_timeout_cb(void* arg); +static void cm_exchange_timeout_cb(void* arg); +static void cm_candidate_ping_timer_cb(void* arg); + +struct cm_reverse_pending { + struct cm_reverse_pending* next; + uint32_t request_id; + struct CONN_MGR_ENTRY* entry; + struct ETCP_CONN** conns; + uint8_t conn_count; +}; + +struct cm_exchange_pending { + struct cm_exchange_pending* next; + uint32_t request_id; + struct CONN_MGR_ENTRY* entry; + void* timeout_timer; + struct CONN_MGR_INTERM_EXCHANGE_RESP cached_resp; + uint8_t resp_received; + uint8_t probes_done; +}; + +struct CONN_MGR* conn_mgr_init(struct UTUN_INSTANCE* instance) { + if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_init: NULL instance"); return NULL; } + struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR)); + if (!mgr) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_init: alloc failed"); return NULL; } + mgr->instance = instance; + mgr->entry_capacity = 8; + mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY)); + if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_init: entries alloc failed"); u_free(mgr); return NULL; } + etcp_router_bind(instance, ETCP_ID_CONN_MGR, conn_mgr_router_recv_handler); + mgr->bg_ping_timer = uasync_set_timeout(instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); + mgr->candidate_ping_timer = uasync_set_timeout(instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); + mgr->bg_ping_cycle_start_tb = get_time_tb(); + mgr->initialized = 1; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: initialized, bg_ping started"); + return mgr; +} + +void conn_mgr_destroy(struct CONN_MGR* mgr) { + if (!mgr) return; + if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; } + if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; } + etcp_router_unbind(mgr->instance, ETCP_ID_CONN_MGR); + for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]); + u_free(mgr->entries); + mgr->initialized = 0; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: destroyed, entries=%zu", (size_t)mgr->entry_count); + u_free(mgr); +} + +static void cm_entry_destroy(struct CONN_MGR_ENTRY* entry) { + if (!entry || entry->state == CONN_MGR_STATE_DISCONNECTED) return; + if (entry->idle_timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->idle_timer); entry->idle_timer = NULL; } + if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } + entry->state = CONN_MGR_STATE_DISCONNECTED; + entry->conn_type = CONN_TYPE_NONE; +} + +static struct CONN_MGR_ENTRY* cm_ensure_entry(struct CONN_MGR* mgr, uint64_t node_id) { + struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, node_id); + if (e) return e; + if (mgr->entry_count >= mgr->entry_capacity) { + size_t new_cap = mgr->entry_capacity * 2; + struct CONN_MGR_ENTRY* new_e = u_realloc(mgr->entries, new_cap * sizeof(struct CONN_MGR_ENTRY)); + if (!new_e) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: realloc entries failed"); return NULL; } + memset(new_e + mgr->entry_capacity, 0, (new_cap - mgr->entry_capacity) * sizeof(struct CONN_MGR_ENTRY)); + mgr->entries = new_e; + mgr->entry_capacity = new_cap; + } + struct CONN_MGR_ENTRY* new_entry = &mgr->entries[mgr->entry_count++]; + memset(new_entry, 0, sizeof(*new_entry)); + new_entry->node_id = node_id; + new_entry->mgr = mgr; + return new_entry; +} + +static struct CONN_MGR_ENTRY* cm_find_entry(struct CONN_MGR* mgr, uint64_t node_id) { + for (size_t i = 0; i < mgr->entry_count; i++) + if (mgr->entries[i].node_id == node_id && mgr->entries[i].state != CONN_MGR_STATE_DISCONNECTED) + return &mgr->entries[i]; + return NULL; +} + +static void cm_deliver_result(struct CONN_MGR_ENTRY* entry, int result) { + struct cm_cb_node* node = entry->cb_list; + while (node) { struct cm_cb_node* next = node->next; node->cb(result, entry->node_id, node->arg); u_free(node); node = next; } + entry->cb_list = NULL; +} + +static void cm_update_nodeinfo(struct CONN_MGR_ENTRY* entry) { + struct ROUTE_BGP* bgp = entry->mgr->instance->bgp; + if (!bgp) return; + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, entry->node_id); + if (!nq) return; + nq->conn_mgr_type = entry->conn_type; + memcpy(nq->conn_mgr_intermediaries, entry->intermediaries, sizeof(entry->intermediaries)); + nq->conn_mgr_intermediariy_count = entry->intermediariy_count; +} + +static void cm_clear_nodeinfo(struct CONN_MGR* mgr, uint64_t node_id) { + struct NODEINFO_Q* nq = nodeinfo_find_by_id(mgr->instance->bgp, node_id); + if (!nq) return; + nq->conn_mgr_type = CONN_TYPE_NONE; + nq->conn_mgr_intermediariy_count = 0; +} + +int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms, + conn_mgr_connect_callback_t cb, void* cb_arg) { + if (!mgr) return CONN_MGR_ERR_INTERNAL; + struct ROUTE_BGP* bgp = mgr->instance->bgp; + if (!bgp) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_connect: BGP not initialized"); return CONN_MGR_ERR_INTERNAL; } + struct NODEINFO_Q* target = nodeinfo_find_by_id(bgp, node_id); + if (!target) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_connect: node 0x%016llx not found", (unsigned long long)node_id); return CONN_MGR_ERR_NOT_FOUND; } + struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); + if (entry && entry->state == CONN_MGR_STATE_CONNECTED) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr_connect: node 0x%016llx already connected", (unsigned long long)node_id); + return CONN_MGR_ERR_ALREADY_CONNECTED; + } + if (entry && entry->state == CONN_MGR_STATE_CONNECTING) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr_connect: node 0x%016llx already connecting, adding cb", (unsigned long long)node_id); + struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node)); + if (cn) { cn->cb = cb; cn->arg = cb_arg; cn->next = entry->cb_list; entry->cb_list = cn; } + return CONN_MGR_OK; + } + struct ETCP_CONN* existing = route_bgp_find_conn_for_node(bgp, node_id); + if (existing && existing->peer_node_id == node_id && existing->links) { + struct ETCP_LINK* l = existing->links; + while (l) { if (l->link_status && l->initialized) break; l = l->next; } + if (l) { + entry = cm_ensure_entry(mgr, node_id); + if (!entry) return CONN_MGR_ERR_INTERNAL; + entry->state = CONN_MGR_STATE_CONNECTED; + entry->conn_type = CONN_TYPE_DIRECT; + entry->idle_timeout_ms = idle_timeout_ms; + entry->alien = target->alien; + entry->last_traffic_tb = get_time_tb(); + { struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node)); + if (cn) { cn->cb = cb; cn->arg = cb_arg; entry->cb_list = cn; } } + cm_update_nodeinfo(entry); + cm_deliver_result(entry, CONN_MGR_OK); + return CONN_MGR_OK; + } + } + entry = cm_ensure_entry(mgr, node_id); + if (!entry) return CONN_MGR_ERR_INTERNAL; + entry->state = CONN_MGR_STATE_CONNECTING; + entry->alien = target->alien; + { struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node)); + if (cn) { cn->cb = cb; cn->arg = cb_arg; cn->next = entry->cb_list; entry->cb_list = cn; } } + entry->idle_timeout_ms = idle_timeout_ms; + entry->last_traffic_tb = get_time_tb(); + entry->local_scan_state = CM_TRY_NONE; + entry->main_connect_state = CM_TRY_NONE; + entry->main.phase = 0; + entry->main.timer = NULL; + entry->main.request_id = 0; + int has_our_local = cm_has_local_addr(bgp->local_node); + int has_target_local = cm_has_local_addr(target); + if (has_our_local && has_target_local) { entry->local_scan_state = CM_TRY_PENDING; cm_start_local_scan(entry); } + entry->main_connect_state = CM_TRY_PENDING; + cm_start_phase_direct(entry); + return CONN_MGR_OK; +} + +int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id) { + if (!mgr) return CONN_MGR_ERR_INTERNAL; + struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); + if (!entry) return CONN_MGR_ERR_NOT_FOUND; + struct CONN_MGR_DISCONNECT pkt; + pkt.cmd = ETCP_ID_CONN_MGR; + pkt.subcmd = CONN_MGR_SUBCMD_DISCONNECT; + pkt.node_id = node_id; + struct ll_entry* qe = queue_entry_new(sizeof(pkt)); + if (qe) { memcpy(qe->data, &pkt, sizeof(pkt)); etcp_route_send(mgr->instance, node_id, qe, 1); } + cm_clear_nodeinfo(mgr, node_id); + cm_entry_destroy(entry); + entry->state = CONN_MGR_STATE_DISCONNECTED; + entry->conn_type = CONN_TYPE_NONE; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: disconnected node 0x%016llx", (unsigned long long)node_id); + return CONN_MGR_OK; +} + +int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t timeout_ms) { + if (!mgr) return CONN_MGR_ERR_INTERNAL; + struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); + if (!entry) return CONN_MGR_ERR_NOT_FOUND; + entry->idle_timeout_ms = timeout_ms; + if (entry->state == CONN_MGR_STATE_CONNECTED && timeout_ms > 0 && !entry->idle_timer) + entry->idle_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_IDLE_CHECK_INTERVAL_TB, entry, cm_idle_timer_cb, "conn_mgr_idle"); + if (timeout_ms == 0 && entry->idle_timer) { uasync_cancel_timeout(mgr->instance->ua, entry->idle_timer); entry->idle_timer = NULL; } + return CONN_MGR_OK; +} + +int conn_mgr_get_status(struct CONN_MGR* mgr, uint64_t node_id, uint8_t* out_state, uint8_t* out_conn_type) { + if (!mgr || !out_state || !out_conn_type) return CONN_MGR_ERR_INTERNAL; + struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); + if (!entry) { *out_state = CONN_MGR_STATE_DISCONNECTED; *out_conn_type = CONN_TYPE_NONE; return CONN_MGR_OK; } + *out_state = entry->state; + *out_conn_type = entry->conn_type; + return CONN_MGR_OK; +} + +int conn_mgr_add_alien_node(struct CONN_MGR* mgr, const uint8_t* nodeinfo_data, size_t len) { + if (!mgr || !nodeinfo_data || len < sizeof(struct NODEINFO)) return CONN_MGR_ERR_INTERNAL; + struct ROUTE_BGP* bgp = mgr->instance->bgp; + if (!bgp) return CONN_MGR_ERR_INTERNAL; + const struct NODEINFO* ni = (const struct NODEINFO*)nodeinfo_data; + uint64_t node_id = ni->node_id; + struct NODEINFO_Q* existing = nodeinfo_find_by_id(bgp, node_id); + if (existing) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: alien node 0x%016llx already exists, ignoring", (unsigned long long)node_id); return CONN_MGR_OK; } + size_t dyn_sz = len - sizeof(struct NODEINFO); + size_t data_sz = sizeof(struct NODEINFO_Q) - sizeof(struct ll_entry) + dyn_sz + 8; + struct NODEINFO_Q* nq = (struct NODEINFO_Q*)queue_entry_new(data_sz); + if (!nq) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: alien node alloc failed"); return CONN_MGR_ERR_INTERNAL; } + memcpy(&nq->node, nodeinfo_data, len); + nq->alien = 1; + nq->dirty = 0; + nq->last_ver = ni->ver; + nq->connectivity.ping_req_time = 0; + if (bgp->nodes) queue_data_put_with_index(bgp->nodes, &nq->ll); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: added alien node 0x%016llx", (unsigned long long)node_id); + return CONN_MGR_OK; +} + +int conn_mgr_send(struct CONN_MGR* mgr, uint64_t node_id, struct ll_entry* entry) { + if (!mgr || !entry) return -1; + struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, node_id); + if (!e || e->state != CONN_MGR_STATE_CONNECTED) { queue_entry_free(entry); return -1; } + e->last_traffic_tb = get_time_tb(); + struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(mgr->instance, node_id, ETCP_ID_CONN_MGR); + if (!rconn) { queue_entry_free(entry); return -1; } + int ret = etcp_router_conn_send(rconn, entry->dgram, entry->len); + queue_entry_free(entry); + return ret; +} + +void conn_mgr_update_best_candidates(struct CONN_MGR* mgr, uint64_t node_id, uint16_t rtt) { + if (!mgr || node_id == mgr->instance->node_id || rtt == 0) return; + for (uint8_t i = 0; i < mgr->best_candidate_count; i++) { + if (mgr->best_candidates[i].node_id == node_id) { + mgr->best_candidates[i].rtt = rtt; + for (uint8_t j = i; j > 0 && mgr->best_candidates[j].rtt < mgr->best_candidates[j-1].rtt; j--) { + struct CONN_MGR_CANDIDATE tmp = mgr->best_candidates[j]; + mgr->best_candidates[j] = mgr->best_candidates[j-1]; + mgr->best_candidates[j-1] = tmp; + } + for (uint8_t j = i; j + 1 < mgr->best_candidate_count && mgr->best_candidates[j].rtt > mgr->best_candidates[j+1].rtt; j++) { + struct CONN_MGR_CANDIDATE tmp = mgr->best_candidates[j]; + mgr->best_candidates[j] = mgr->best_candidates[j+1]; + mgr->best_candidates[j+1] = tmp; + } + return; + } + } + if (mgr->best_candidate_count < CONN_MGR_MAX_CANDIDATES) { + mgr->best_candidates[mgr->best_candidate_count].node_id = node_id; + mgr->best_candidates[mgr->best_candidate_count].rtt = rtt; + mgr->best_candidate_count++; + } else if (rtt < mgr->best_candidates[CONN_MGR_MAX_CANDIDATES - 1].rtt) { + mgr->best_candidates[CONN_MGR_MAX_CANDIDATES - 1].node_id = node_id; + mgr->best_candidates[CONN_MGR_MAX_CANDIDATES - 1].rtt = rtt; + } else return; + for (uint8_t i = mgr->best_candidate_count - 1; i > 0; i--) { + if (mgr->best_candidates[i].rtt >= mgr->best_candidates[i-1].rtt) break; + struct CONN_MGR_CANDIDATE tmp = mgr->best_candidates[i]; + mgr->best_candidates[i] = mgr->best_candidates[i-1]; + mgr->best_candidates[i-1] = tmp; + } +} + +static void cm_idle_timer_cb(void* arg) { + struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; + struct CONN_MGR* mgr = entry->mgr; + uint64_t now_tb = get_time_tb(); + if (entry->state != CONN_MGR_STATE_CONNECTED || entry->idle_timeout_ms == 0) return; + uint64_t idle = now_tb - entry->last_traffic_tb; + if (idle > (uint64_t)entry->idle_timeout_ms * 10) { + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: node 0x%016llx idle %llums, disconnecting", + (unsigned long long)entry->node_id, (unsigned long long)(idle / 10)); + conn_mgr_disconnect_node(mgr, entry->node_id); + return; + } + entry->idle_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_IDLE_CHECK_INTERVAL_TB, entry, cm_idle_timer_cb, "conn_mgr_idle"); +} + +static int cm_nat_compatible(struct ETCP_SOCKET* our, uint8_t target_meta_type) { + uint8_t tt = target_meta_type; + if (tt == NAT_VERIFIED_DIRECT) tt = CFG_SERVER_TYPE_PUBLIC; + else if (tt >= NAT_VERIFIED_UNKNOWN) tt = CFG_SERVER_TYPE_PUBLIC; + uint8_t ot = our->type; + if (ot == CFG_SERVER_TYPE_PUBLIC) return 1; + if (ot == CFG_SERVER_TYPE_NAT && our->nat_type == NAT_TYPE_EIM && tt == CFG_SERVER_TYPE_PUBLIC) return 1; + if (ot == CFG_SERVER_TYPE_NAT && our->nat_type == NAT_TYPE_EIM && tt == CFG_SERVER_TYPE_NAT) return 1; + if (ot == CFG_SERVER_TYPE_PRIVATE && tt == CFG_SERVER_TYPE_PRIVATE) return 2; + return 0; +} + +static uint16_t cm_get_node_min_rtt(struct NODEINFO_Q* nq) { + uint16_t rtt = 0xFFFF; + if (nq->connectivity.interface_status == PROBE_RESULT_REACHABLE && nq->connectivity.interface_min_rtt < rtt) + rtt = nq->connectivity.interface_min_rtt; + if (nq->connectivity.nat_status == PROBE_RESULT_REACHABLE && nq->connectivity.nat_min_rtt < rtt) + rtt = nq->connectivity.nat_min_rtt; + if (nq->connectivity.real_status == PROBE_RESULT_REACHABLE && nq->connectivity.real_min_rtt < rtt) + rtt = nq->connectivity.real_min_rtt; + return rtt; +} + +static uint64_t cm_get_node_max_probe_time(struct NODEINFO_Q* nq) { + uint64_t t = nq->connectivity.interface_probe_time; + if (nq->connectivity.nat_probe_time > t) t = nq->connectivity.nat_probe_time; + if (nq->connectivity.real_probe_time > t) t = nq->connectivity.real_probe_time; + return t; +} + +static int cm_is_rtt_fresh(struct NODEINFO_Q* nq, uint64_t now_tb) { + uint64_t last = cm_get_node_max_probe_time(nq); + return (last > 0 && (now_tb - last) < (uint64_t)(CONN_MGR_CANDIDATE_CACHE_MS * 10)); +} + +static int cm_has_direct_ip(struct NODEINFO_Q* nq) { + if (!nq) return 0; + const struct NODEINFO_IPV4_ADDR* addrs = NULL; + const struct NODEINFO_IPV4_SOCKET_META* metas = NULL; + int addr_count = get_node_v4_addrs(nq, &addrs); + int meta_count = get_node_v4_sockets_meta(nq, &metas); + for (int i = 0; i < addr_count; i++) { + if (addrs[i].type == ADDR_TYPE_REAL || addrs[i].type == ADDR_TYPE_NAT) { + for (int j = 0; j < meta_count; j++) { + if (metas[j].id == addrs[i].socket_id && + (metas[j].type == CFG_SERVER_TYPE_PUBLIC || metas[j].type == NAT_VERIFIED_EIM || metas[j].type == NAT_VERIFIED_DIRECT)) + return 1; + } + } + } + return 0; +} + +static int cm_has_local_addr(struct NODEINFO_Q* nq) { + if (!nq) return 0; + const struct NODEINFO_IPV4_ADDR* addrs = NULL; + const struct NODEINFO_IPV4_SOCKET_META* metas = NULL; + int addr_count = get_node_v4_addrs(nq, &addrs); + int meta_count = get_node_v4_sockets_meta(nq, &metas); + for (int i = 0; i < addr_count; i++) { + if (addrs[i].type == ADDR_TYPE_INTERFACE) { + for (int j = 0; j < meta_count; j++) { + if (metas[j].id == addrs[i].socket_id && metas[j].type == CFG_SERVER_TYPE_PRIVATE) + return 1; + } + } + } + return 0; +} + +static void cm_ping_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, + const uint8_t* resp_data, size_t resp_data_len); +static void cm_direct_timeout_cb(void* arg); + +struct cm_ping_ctx { + struct CONN_MGR_ENTRY* entry; + struct sockaddr_storage addr; + struct ETCP_SOCKET* sock; + uint8_t phase; /* 0=local_scan, 1=direct */ + uint8_t attempt; +}; + +static void cm_start_local_scan(struct CONN_MGR_ENTRY* entry) { + struct ROUTE_BGP* bgp = entry->mgr->instance->bgp; + struct NODEINFO_Q* target = nodeinfo_find_by_id(bgp, entry->node_id); + if (!target) { entry->local_scan_state = CM_TRY_FAILED; return; } + const struct NODEINFO_IPV4_ADDR* addrs = NULL; + const struct NODEINFO_IPV4_SOCKET_META* metas = NULL; + int addr_count = get_node_v4_addrs(target, &addrs); + int meta_count = get_node_v4_sockets_meta(target, &metas); + for (int i = 0; i < addr_count; i++) { + if (addrs[i].type != ADDR_TYPE_INTERFACE) continue; + for (int j = 0; j < meta_count; j++) { + if (metas[j].id != addrs[i].socket_id || metas[j].type != CFG_SERVER_TYPE_PRIVATE) continue; + struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets; + while (s) { + if (s->type == CFG_SERVER_TYPE_PRIVATE && s->local_addr.ss_family == AF_INET) { + struct sockaddr_in sin; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); + sin.sin_port = htons(addrs[i].port); + struct sockaddr_storage sa; + memset(&sa, 0, sizeof(sa)); + memcpy(&sa, &sin, sizeof(sin)); + struct cm_ping_ctx* ctx = u_calloc(1, sizeof(struct cm_ping_ctx)); + if (!ctx) continue; + ctx->entry = entry; + ctx->addr = sa; + ctx->sock = s; + ctx->phase = 0; + ctx->attempt = 0; + etcp_send_ping_to_socket(entry->mgr->instance, s, target->node.public_key, &sa, + CONN_MGR_LOCAL_SCAN_TIMEOUT_MS, cm_ping_cb, ctx, NULL, 0); + return; + } + s = s->next; + } + } + } + entry->local_scan_state = CM_TRY_FAILED; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: local scan no candidates for 0x%016llx", (unsigned long long)entry->node_id); +} + +static void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) { + if (entry->main_connect_state == CM_TRY_OK) return; + if (entry->local_scan_state == CM_TRY_OK) return; + struct ROUTE_BGP* bgp = entry->mgr->instance->bgp; + struct NODEINFO_Q* target = nodeinfo_find_by_id(bgp, entry->node_id); + if (!target) { cm_deliver_result(entry, CONN_MGR_ERR_NOT_FOUND); return; } + const struct NODEINFO_IPV4_ADDR* addrs = NULL; + const struct NODEINFO_IPV4_SOCKET_META* metas = NULL; + int addr_count = get_node_v4_addrs(target, &addrs); + int meta_count = get_node_v4_sockets_meta(target, &metas); + uint8_t priority_order[] = { ADDR_TYPE_REAL, ADDR_TYPE_NAT, ADDR_TYPE_INTERFACE }; + for (int pri = 0; pri < 3; pri++) { + for (int i = 0; i < addr_count; i++) { + if (addrs[i].type != priority_order[pri]) continue; + for (int j = 0; j < meta_count; j++) { + if (metas[j].id != addrs[i].socket_id) continue; + struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets; + while (s) { + int compat = cm_nat_compatible(s, metas[j].type); + if (compat && s->local_addr.ss_family == AF_INET) { + struct sockaddr_in sin; + memset(&sin, 0, sizeof(sin)); + sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); + sin.sin_port = htons(addrs[i].port); + struct sockaddr_storage sa; + memset(&sa, 0, sizeof(sa)); + memcpy(&sa, &sin, sizeof(sin)); + struct ETCP_CONN* conn = etcp_connection_create(entry->mgr->instance, NULL); + if (conn) { + etcp_conn_set_ready_cbk(conn, cm_direct_ready_cb, entry); + sc_set_peer_public_key(&conn->crypto_ctx, target->node.public_key, 0); + struct ETCP_LINK* link = etcp_link_new(conn, s, &sa, 0); + if (link) { + entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, + CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS * 10, entry, cm_direct_timeout_cb, "conn_mgr_direct"); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 1 (direct) started INIT to 0x%016llx", + (unsigned long long)entry->node_id); + return; + } + } + } + s = s->next; + } + } + } + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 1 (direct) no compatible pairs for 0x%016llx", + (unsigned long long)entry->node_id); + int has_our = cm_has_direct_ip(bgp->local_node); + int has_target = cm_has_direct_ip(target); + if (has_our && !has_target) { cm_start_phase_reverse(entry); return; } + cm_start_phase_indirect(entry); +} + +static void cm_direct_timeout_cb(void* arg) { + struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; + entry->main.timer = NULL; + struct ROUTE_BGP* bgp = entry->mgr->instance->bgp; + struct NODEINFO_Q* target = nodeinfo_find_by_id(bgp, entry->node_id); + int has_our = cm_has_direct_ip(bgp->local_node); + int has_target = target ? cm_has_direct_ip(target) : 0; + if (entry->local_scan_state != CM_TRY_OK) { + if (has_our && !has_target) { cm_start_phase_reverse(entry); return; } + cm_start_phase_indirect(entry); + return; + } + if (entry->local_scan_state == CM_TRY_FAILED && entry->main_connect_state == CM_TRY_FAILED) + cm_deliver_result(entry, CONN_MGR_ERR_UNREACHABLE); +} + +static void cm_direct_ready_cb(struct ETCP_CONN* conn, void* arg) { + struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; + if (!conn || conn->peer_node_id != entry->node_id) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: direct ready cb peer mismatch conn=%p peer=0x%llx entry=0x%llx", + (void*)conn, (unsigned long long)conn->peer_node_id, (unsigned long long)entry->node_id); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: direct ready for 0x%016llx", (unsigned long long)entry->node_id); + if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } + entry->main_connect_state = CM_TRY_OK; + entry->conn_type = CONN_TYPE_DIRECT; + entry->state = CONN_MGR_STATE_CONNECTED; + cm_update_nodeinfo(entry); + cm_deliver_result(entry, CONN_MGR_OK); +} + +static void cm_ping_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, + const uint8_t* resp_data, size_t resp_data_len) { + struct cm_ping_ctx* ctx = (struct cm_ping_ctx*)arg; + struct CONN_MGR_ENTRY* entry = ctx->entry; + (void)nonce; (void)resp_data; (void)resp_data_len; + if (entry->main_connect_state == CM_TRY_OK) { u_free(ctx); return; } + if (success && ctx->phase == 0) { + entry->local_scan_state = CM_TRY_OK; + entry->conn_type = CONN_TYPE_DIRECT; + entry->state = CONN_MGR_STATE_CONNECTED; + if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: local scan OK for 0x%016llx rtt=%u", (unsigned long long)entry->node_id, rtt); + cm_update_nodeinfo(entry); + cm_deliver_result(entry, CONN_MGR_OK); + u_free(ctx); + return; + } + ctx->attempt++; + if (ctx->attempt < CONN_MGR_LOCAL_SCAN_ATTEMPTS && ctx->phase == 0) { + etcp_send_ping_to_socket(entry->mgr->instance, ctx->sock, + nodeinfo_find_by_id(entry->mgr->instance->bgp, entry->node_id)->node.public_key, + &ctx->addr, CONN_MGR_LOCAL_SCAN_TIMEOUT_MS, cm_ping_cb, ctx, NULL, 0); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: %s failed for 0x%016llx after %d attempts", + ctx->phase == 0 ? "local scan" : "direct", (unsigned long long)entry->node_id, ctx->attempt); + if (ctx->phase == 0) entry->local_scan_state = CM_TRY_FAILED; + u_free(ctx); +} + +static void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) { + struct ROUTE_BGP* bgp = entry->mgr->instance->bgp; + struct NODEINFO_Q* local = bgp->local_node; + if (!local) { cm_deliver_result(entry, CONN_MGR_ERR_INTERNAL); return; } + uint32_t req_id = ++entry->mgr->next_request_id; + entry->main.request_id = req_id; + entry->main.phase = 2; + const struct NODEINFO_IPV4_ADDR* addrs = NULL; + const struct NODEINFO_IPV4_SOCKET_META* metas = NULL; + int addr_count = get_node_v4_addrs(local, &addrs); + int meta_count = get_node_v4_sockets_meta(local, &metas); + uint8_t direct_count = 0; + struct { uint8_t type; uint8_t ip[4]; uint16_t port; uint8_t socket_id; } out_addrs[8]; + for (int i = 0; i < addr_count && direct_count < 8; i++) { + if (addrs[i].type != ADDR_TYPE_REAL && addrs[i].type != ADDR_TYPE_NAT) continue; + for (int j = 0; j < meta_count; j++) { + if (metas[j].id == addrs[i].socket_id && + (metas[j].type == CFG_SERVER_TYPE_PUBLIC || metas[j].type == NAT_VERIFIED_EIM || metas[j].type == NAT_VERIFIED_DIRECT)) { + out_addrs[direct_count].type = addrs[i].type; + memcpy(out_addrs[direct_count].ip, addrs[i].addr, 4); + out_addrs[direct_count].port = htons(addrs[i].port); + out_addrs[direct_count].socket_id = addrs[i].socket_id; + direct_count++; + break; + } + } + } + if (direct_count == 0) { cm_start_phase_indirect(entry); 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) { cm_deliver_result(entry, CONN_MGR_ERR_INTERNAL); return; } + struct CONN_MGR_DIRECT_REQ* req = (struct CONN_MGR_DIRECT_REQ*)pkt; + req->cmd = ETCP_ID_CONN_MGR; + req->subcmd = CONN_MGR_SUBCMD_DIRECT_REQ; + req->request_id = req_id; + 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(pkt_size); + if (!qe) { u_free(pkt); cm_deliver_result(entry, CONN_MGR_ERR_INTERNAL); return; } + memcpy(qe->data, pkt, pkt_size); + qe->len = (uint16_t)pkt_size; + u_free(pkt); + etcp_route_send(entry->mgr->instance, entry->node_id, qe, 1); + + struct cm_reverse_pending* rp = u_calloc(1, sizeof(struct cm_reverse_pending)); + if (rp) { rp->request_id = req_id; rp->entry = entry; rp->next = entry->mgr->reverse_pending; + entry->mgr->reverse_pending = rp; } + entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS * 10, + entry, cm_reverse_timeout_cb, "conn_mgr_reverse"); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 2 (reverse) sent DIRECT_REQ to 0x%016llx with %u addrs", + (unsigned long long)entry->node_id, direct_count); +} + +static void cm_reverse_timeout_cb(void* arg) { + struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; + entry->main.timer = NULL; + struct cm_reverse_pending** pp = &entry->mgr->reverse_pending; + while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); break; } pp = &(*pp)->next; } + cm_start_phase_indirect(entry); +} + +static void cm_reverse_ready_cb(struct ETCP_CONN* conn, void* arg) { + struct cm_reverse_pending* rp = (struct cm_reverse_pending*)arg; + struct CONN_MGR_ENTRY* entry = rp->entry; + if (!conn || conn->peer_node_id != entry->node_id) { u_free(rp); return; } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse ready for 0x%016llx", (unsigned long long)entry->node_id); + if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } + struct cm_reverse_pending** pp = &entry->mgr->reverse_pending; + while (*pp) { if (*pp == rp) { *pp = rp->next; break; } pp = &(*pp)->next; } + entry->main_connect_state = CM_TRY_OK; + entry->conn_type = CONN_TYPE_REVERSE; + entry->state = CONN_MGR_STATE_CONNECTED; + cm_update_nodeinfo(entry); + cm_deliver_result(entry, CONN_MGR_OK); + u_free(rp); +} + +static int cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry) { + struct CONN_MGR* mgr = entry->mgr; + if (mgr->best_candidate_count == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 3 (indirect) no candidates for 0x%016llx", + (unsigned long long)entry->node_id); + return CONN_MGR_ERR_UNREACHABLE; + } + uint32_t req_id = ++mgr->next_request_id; + entry->main.request_id = req_id; + entry->main.phase = 3; + struct CONN_MGR_INTERM_EXCHANGE_REQ req; + memset(&req, 0, sizeof(req)); + req.cmd = ETCP_ID_CONN_MGR; + req.subcmd = CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ; + req.request_id = req_id; + req.candidate_count = mgr->best_candidate_count > 4 ? 4 : mgr->best_candidate_count; + for (uint8_t i = 0; i < req.candidate_count; i++) + req.candidates[i] = mgr->best_candidates[i]; + struct ll_entry* qe = queue_entry_new(sizeof(req)); + if (!qe) return CONN_MGR_ERR_INTERNAL; + memcpy(qe->data, &req, sizeof(req)); + etcp_route_send(entry->mgr->instance, entry->node_id, qe, 1); + entry->main.timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10, + entry, cm_exchange_timeout_cb, "conn_mgr_interm_exch"); + { + struct cm_exchange_pending* ep = u_calloc(1, sizeof(struct cm_exchange_pending)); + if (ep) { ep->request_id = req_id; ep->entry = entry; ep->timeout_timer = entry->main.timer; + ep->next = mgr->exchange_pending; mgr->exchange_pending = ep; } + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 3 (indirect) sent EXCHANGE_REQ to 0x%016llx with %u candidates", + (unsigned long long)entry->node_id, req.candidate_count); + return CONN_MGR_OK; +} + +static void cm_exchange_timeout_cb(void* arg) { + struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; + entry->main.timer = NULL; + struct cm_exchange_pending** pp = &entry->mgr->exchange_pending; + while (*pp) { if ((*pp)->entry == entry) { struct cm_exchange_pending* ep = *pp; *pp = ep->next; u_free(ep); break; } pp = &(*pp)->next; } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: exchange timeout for 0x%016llx", (unsigned long long)entry->node_id); + cm_deliver_result(entry, CONN_MGR_ERR_TIMEOUT); +} + +static void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CONN_MGR_INTERM_EXCHANGE_RESP* resp) { + struct { uint64_t node_id; uint64_t total_rtt; } all[8]; + uint8_t count = 0; + + for (uint8_t i = 0; i < resp->your_count && count < 8; i++) { + uint64_t id = resp->your_candidates[i].node_id; + uint16_t our_rtt = 0; + for (uint8_t j = 0; j < CONN_MGR_MAX_CANDIDATES; j++) { + if (entry->mgr->best_candidates[j].node_id == id) { our_rtt = entry->mgr->best_candidates[j].rtt; break; } + } + all[count].node_id = id; + all[count].total_rtt = (uint64_t)our_rtt + resp->your_candidates[i].rtt; + count++; + } + for (uint8_t i = 0; i < resp->my_count && count < 8; i++) { + uint64_t id = resp->my_candidates[i].node_id; + struct NODEINFO_Q* nq = nodeinfo_find_by_id(entry->mgr->instance->bgp, id); + uint16_t our_rtt = nq ? cm_get_node_min_rtt(nq) : 0; + if (our_rtt == 0xFFFF) our_rtt = 0; + all[count].node_id = id; + all[count].total_rtt = (uint64_t)our_rtt + resp->my_candidates[i].rtt; + count++; + } + + for (uint8_t i = 0; i < count; i++) + for (uint8_t j = i + 1; j < count; j++) + if (all[j].total_rtt < all[i].total_rtt) { typeof(all[0]) t = all[i]; all[i] = all[j]; all[j] = t; } + + uint8_t sel = count > CONN_MGR_MAX_INTERMEDIARIES ? CONN_MGR_MAX_INTERMEDIARIES : count; + for (uint8_t i = 0; i < sel; i++) entry->intermediaries[i] = all[i].node_id; + entry->intermediariy_count = sel; + entry->rr_idx = 0; + + struct CONN_MGR_INTERM_SELECTED pkt; + memset(&pkt, 0, sizeof(pkt)); + pkt.cmd = ETCP_ID_CONN_MGR; + pkt.subcmd = CONN_MGR_SUBCMD_INTERM_SELECTED; + pkt.request_id = 0; + pkt.count = sel; + for (uint8_t i = 0; i < sel; i++) { + pkt.selected[i].node_id = all[i].node_id; + pkt.selected[i].rtt = 0; + struct NODEINFO_Q* nq = nodeinfo_find_by_id(entry->mgr->instance->bgp, all[i].node_id); + if (nq) pkt.selected[i].rtt = cm_get_node_min_rtt(nq); + } + struct ll_entry* qe = queue_entry_new(sizeof(pkt)); + if (qe) { memcpy(qe->data, &pkt, sizeof(pkt)); etcp_route_send(entry->mgr->instance, entry->node_id, qe, 1); } + + entry->conn_type = CONN_TYPE_INDIRECT; + entry->state = CONN_MGR_STATE_CONNECTED; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: indirect OK for 0x%016llx via %u intermediaries", + (unsigned long long)entry->node_id, sel); + cm_update_nodeinfo(entry); + cm_deliver_result(entry, CONN_MGR_OK); +} + +static void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { + const struct CONN_MGR_INTERM_EXCHANGE_RESP* resp = (const struct CONN_MGR_INTERM_EXCHANGE_RESP*)data; + size_t min_size = offsetof(struct CONN_MGR_INTERM_EXCHANGE_RESP, your_candidates); + if (len < min_size) return; + struct cm_exchange_pending* ep = mgr->exchange_pending; + while (ep) { if (ep->request_id == resp->request_id) break; ep = ep->next; } + if (!ep) return; + + if (ep->entry->main.timer) { uasync_cancel_timeout(mgr->instance->ua, ep->entry->main.timer); ep->entry->main.timer = NULL; } + + memcpy(&ep->cached_resp, resp, len < sizeof(ep->cached_resp) ? len : sizeof(ep->cached_resp)); + ep->resp_received = 1; + ep->probes_done = 0; + + uint8_t need_probe = 0; + uint64_t now_tb = get_time_tb(); + for (uint8_t i = 0; i < ep->cached_resp.my_count && i < 4; i++) { + struct NODEINFO_Q* nq = nodeinfo_find_by_id(mgr->instance->bgp, ep->cached_resp.my_candidates[i].node_id); + if (!nq || !cm_is_rtt_fresh(nq, now_tb)) { need_probe++; route_connectivity_probe_node(mgr->instance, nq); } + } + if (need_probe == 0) { ep->probes_done = 1; cm_compute_intermediaries(ep->entry, &ep->cached_resp); + struct cm_exchange_pending** pp = &mgr->exchange_pending; + while (*pp) { if (*pp == ep) { *pp = ep->next; break; } pp = &(*pp)->next; } + u_free(ep); return; } + uasync_set_timeout(mgr->instance->ua, 2000, ep, NULL, "cm_exch_probe"); +} + +static void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { + const struct CONN_MGR_INTERM_SELECTED* sel = (const struct CONN_MGR_INTERM_SELECTED*)data; + if (len < offsetof(struct CONN_MGR_INTERM_SELECTED, selected)) return; + struct CONN_MGR_ENTRY* entry = NULL; + for (size_t i = 0; i < mgr->entry_count; i++) + if (mgr->entries[i].state == CONN_MGR_STATE_CONNECTING) { entry = &mgr->entries[i]; break; } + if (!entry) return; + uint8_t cnt = sel->count > CONN_MGR_MAX_INTERMEDIARIES ? CONN_MGR_MAX_INTERMEDIARIES : sel->count; + for (uint8_t i = 0; i < cnt; i++) entry->intermediaries[i] = sel->selected[i].node_id; + entry->intermediariy_count = cnt; + entry->rr_idx = 0; + entry->conn_type = CONN_TYPE_INDIRECT; + entry->state = CONN_MGR_STATE_CONNECTED; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: got INTERM_SELECTED count=%u", cnt); +} + +static void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!entry || !entry->dgram || entry->len < 3) return; + uint8_t* dgram = entry->dgram + 1; // skip svc_id byte from router + size_t len = entry->len - 1; + uint8_t subcmd = dgram[1]; + struct CONN_MGR* mgr = conn->instance->conn_mgr; + + if (mgr) { + struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, conn->peer_node_id); + if (e && e->state == CONN_MGR_STATE_CONNECTED) e->last_traffic_tb = get_time_tb(); + } + + switch (subcmd) { + case CONN_MGR_SUBCMD_DIRECT_REQ: + cm_handle_direct_req(conn, dgram, len); + break; + case CONN_MGR_SUBCMD_DIRECT_RESP: + break; + case CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ: + if (len >= CONN_MGR_INTERM_EXCHANGE_REQ_SIZE) + cm_handle_interm_exchange_req(conn, (struct CONN_MGR_INTERM_EXCHANGE_REQ*)dgram); + break; + case CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP: + if (mgr) cm_handle_interm_exchange_resp(mgr, dgram, len); + break; + case CONN_MGR_SUBCMD_INTERM_SELECTED: + if (mgr) cm_handle_interm_selected(mgr, dgram, len); + break; + case CONN_MGR_SUBCMD_DISCONNECT: + if (len >= CONN_MGR_DISCONNECT_SIZE && mgr) { + struct CONN_MGR_DISCONNECT* disc = (struct CONN_MGR_DISCONNECT*)dgram; + cm_handle_disconnect(mgr, disc->node_id); + } + break; + } +} + +static void cm_handle_direct_req(struct ETCP_CONN* conn, const uint8_t* data, size_t len) { + struct CONN_MGR* mgr = conn->instance->conn_mgr; + if (!mgr) return; + struct CONN_MGR_DIRECT_REQ* req = (struct CONN_MGR_DIRECT_REQ*)data; + size_t expected = CONN_MGR_DIRECT_REQ_HDR_SIZE + (size_t)req->addr_count * 8; + if (len < expected) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: DIRECT_REQ too short"); return; } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: got DIRECT_REQ from 0x%016llx with %u addrs", + (unsigned long long)conn->peer_node_id, req->addr_count); + uint8_t* p = (uint8_t*)data + CONN_MGR_DIRECT_REQ_HDR_SIZE; + uint64_t src_node_id = conn->peer_node_id; + for (uint8_t i = 0; i < req->addr_count; i++) { + uint8_t type = *p++; + uint8_t ip[4]; memcpy(ip, p, 4); p += 4; + uint16_t port; memcpy(&port, p, 2); p += 2; + uint8_t sock_id = *p++; + (void)sock_id; (void)type; + struct ETCP_SOCKET* s = conn->instance->etcp_sockets; + while (s) { + if (s->type == CFG_SERVER_TYPE_PUBLIC && s->local_addr.ss_family == AF_INET) { + struct sockaddr_in sin; + memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, ip, 4); sin.sin_port = port; + struct sockaddr_storage sa; + memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); + struct ETCP_CONN* new_conn = etcp_connection_create(conn->instance, NULL); + if (new_conn) { + struct NODEINFO_Q* nq = nodeinfo_find_by_id(conn->instance->bgp, src_node_id); + if (nq) sc_set_peer_public_key(&new_conn->crypto_ctx, nq->node.public_key, 0); + struct cm_reverse_pending* rp = u_calloc(1, sizeof(struct cm_reverse_pending)); + if (rp) { rp->request_id = req->request_id; rp->entry = NULL; rp->next = mgr->reverse_pending; + etcp_conn_set_ready_cbk(new_conn, cm_reverse_ready_cb, rp); mgr->reverse_pending = rp; } + etcp_link_new(new_conn, s, &sa, 0); + } + break; + } + s = s->next; + } + } +} + +static void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, + struct CONN_MGR_INTERM_EXCHANGE_REQ* req) { + struct CONN_MGR* mgr = conn->instance->conn_mgr; + if (!mgr) return; + uint64_t now_tb = get_time_tb(); + for (uint8_t i = 0; i < req->candidate_count && i < 4; i++) { + uint64_t cand_id = req->candidates[i].node_id; + struct NODEINFO_Q* nq = nodeinfo_find_by_id(mgr->instance->bgp, cand_id); + if (nq && cm_is_rtt_fresh(nq, now_tb)) continue; + route_connectivity_probe_node(mgr->instance, nq); + } + struct CONN_MGR_INTERM_EXCHANGE_RESP resp; + memset(&resp, 0, sizeof(resp)); + resp.cmd = ETCP_ID_CONN_MGR; + resp.subcmd = CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP; + resp.request_id = req->request_id; + resp.my_count = mgr->best_candidate_count > 4 ? 4 : mgr->best_candidate_count; + for (uint8_t i = 0; i < resp.my_count; i++) + resp.my_candidates[i] = mgr->best_candidates[i]; + resp.your_count = req->candidate_count > 4 ? 4 : req->candidate_count; + for (uint8_t i = 0; i < resp.your_count; i++) { + resp.your_candidates[i].node_id = req->candidates[i].node_id; + struct NODEINFO_Q* nq2 = nodeinfo_find_by_id(mgr->instance->bgp, req->candidates[i].node_id); + uint16_t rtt = nq2 ? cm_get_node_min_rtt(nq2) : 0; + resp.your_candidates[i].rtt = rtt; + } + size_t resp_size = offsetof(struct CONN_MGR_INTERM_EXCHANGE_RESP, your_candidates) + + (size_t)resp.your_count * sizeof(struct CONN_MGR_CANDIDATE); + struct ll_entry* qe = queue_entry_new(resp_size); + if (qe) { + memcpy(qe->data, &resp, resp_size); + etcp_route_send(mgr->instance, conn->peer_node_id, qe, 1); + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent EXCHANGE_RESP to 0x%016llx my=%u your=%u", + (unsigned long long)conn->peer_node_id, resp.my_count, resp.your_count); +} + +static void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id) { + struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); + if (!entry) return; + cm_clear_nodeinfo(mgr, node_id); + cm_entry_destroy(entry); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: remote disconnect from 0x%016llx", (unsigned long long)node_id); +} + +static void cm_bg_ping_timer_cb(void* arg) { + struct CONN_MGR* mgr = (struct CONN_MGR*)arg; + struct ROUTE_BGP* bgp = mgr->instance->bgp; + if (!bgp || !bgp->nodes) { + mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); + return; + } + uint64_t now_tb = get_time_tb(); + struct ll_entry* e = bgp->nodes->head; + size_t total = 0; + while (e) { total++; e = e->next; } + if (total == 0) { + mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); + return; + } + if (mgr->bg_ping_cursor >= total) { + mgr->bg_ping_cursor = 0; + if ((now_tb - mgr->bg_ping_cycle_start_tb) < CONN_MGR_BG_PING_CYCLE_MIN_TB) { + uint64_t delay = CONN_MGR_BG_PING_CYCLE_MIN_TB - (now_tb - mgr->bg_ping_cycle_start_tb); + mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, (int)delay, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); + return; + } + mgr->bg_ping_cycle_start_tb = now_tb; + } + size_t cursor = 0; + e = bgp->nodes->head; + while (e && cursor < mgr->bg_ping_cursor) { cursor++; e = e->next; } + if (e) { + struct NODEINFO_Q* nq = (struct NODEINFO_Q*)e; + uint64_t node_id = nq->node.node_id; + if (node_id != mgr->instance->node_id) { + struct ETCP_CONN* dconn = route_bgp_find_conn_for_node(bgp, node_id); + int has_active = 0; + if (dconn && dconn->links) { + struct ETCP_LINK* l = dconn->links; + while (l) { if (l->initialized && l->link_status) { has_active = 1; break; } l = l->next; } + } + if (!has_active) { + etcp_send_ping(mgr->instance, nq->node.public_key, NULL, CONN_PROBE_TIMEOUT_MS, + NULL, NULL, NULL, 0); + } + } + mgr->bg_ping_cursor++; + } + mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); +} + +static void cm_candidate_ping_timer_cb(void* arg) { + struct CONN_MGR* mgr = (struct CONN_MGR*)arg; + uint64_t now_tb = get_time_tb(); + for (uint8_t i = 0; i < mgr->best_candidate_count; i++) { + uint64_t cand_id = mgr->best_candidates[i].node_id; + struct NODEINFO_Q* nq = nodeinfo_find_by_id(mgr->instance->bgp, cand_id); + if (!nq) { + memmove(&mgr->best_candidates[i], &mgr->best_candidates[i+1], (mgr->best_candidate_count - i - 1) * sizeof(mgr->best_candidates[0])); + mgr->best_candidate_count--; i--; continue; + } + uint64_t last = cm_get_node_max_probe_time(nq); + if (last > 0 && (now_tb - last) >= CONN_MGR_CANDIDATE_STALE_TB) { + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: candidate 0x%016llx stale, removing", (unsigned long long)cand_id); + memmove(&mgr->best_candidates[i], &mgr->best_candidates[i+1], (mgr->best_candidate_count - i - 1) * sizeof(mgr->best_candidates[0])); + mgr->best_candidate_count--; i--; continue; + } + struct ETCP_CONN* dconn = route_bgp_find_conn_for_node(mgr->instance->bgp, cand_id); + int has_active = 0; + if (dconn && dconn->links) { + struct ETCP_LINK* l = dconn->links; + while (l) { if (l->initialized && l->link_status) { has_active = 1; break; } l = l->next; } + } + if (has_active) + etcp_send_ping(mgr->instance, nq->node.public_key, &dconn->links->remote_addr, CONN_PROBE_TIMEOUT_MS, NULL, NULL, NULL, 0); + else + route_connectivity_probe_node(mgr->instance, nq); + } + mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, + mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); +} diff --git a/src/conn_mgr.h b/src/conn_mgr.h new file mode 100644 index 00000000..0ace6bb1 --- /dev/null +++ b/src/conn_mgr.h @@ -0,0 +1,190 @@ +#ifndef CONN_MGR_H +#define CONN_MGR_H + +#include +#include +#include "../lib/ll_queue.h" +#include "etcp_api.h" +#include "etcp_router.h" +#include "route_node.h" + +struct UTUN_INSTANCE; +struct ETCP_CONN; +struct NODEINFO_Q; +struct ROUTE_BGP; + +#define CONN_MGR_MAX_CANDIDATES 3 +#define CONN_MGR_CANDIDATE_CACHE_MS 20000 +#define CONN_MGR_CANDIDATE_STALE_TB 300000 +#define CONN_MGR_CANDIDATE_PING_TB 20000 +#define CONN_MGR_BG_PING_INTERVAL_TB 1000 +#define CONN_MGR_BG_PING_CYCLE_MIN_TB 100000 +#define CONN_MGR_IDLE_CHECK_INTERVAL_TB 10000 +#define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3 +#define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100 +#define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 5000 +#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000 +#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000 + +#define CONN_MGR_OK 0 +#define CONN_MGR_ERR_NOT_FOUND -1 +#define CONN_MGR_ERR_NO_ADDRESSES -2 +#define CONN_MGR_ERR_TIMEOUT -3 +#define CONN_MGR_ERR_UNREACHABLE -4 +#define CONN_MGR_ERR_REFUSED -5 +#define CONN_MGR_ERR_ALREADY_CONNECTED -6 +#define CONN_MGR_ERR_INTERNAL -7 + +enum CONN_MGR_STATE { + CONN_MGR_STATE_DISCONNECTED = 0, + CONN_MGR_STATE_CONNECTING = 1, + CONN_MGR_STATE_CONNECTED = 2, +}; + +enum CM_TRY_STATE { + CM_TRY_NONE = 0, + CM_TRY_PENDING = 1, + CM_TRY_OK = 2, + CM_TRY_FAILED = 3, +}; + +#define CONN_MGR_SUBCMD_DIRECT_REQ 0x01 +#define CONN_MGR_SUBCMD_DIRECT_RESP 0x02 +#define CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ 0x03 +#define CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP 0x04 +#define CONN_MGR_SUBCMD_INTERM_SELECTED 0x05 +#define CONN_MGR_SUBCMD_DISCONNECT 0x06 + +struct CONN_MGR_CANDIDATE { + uint64_t node_id; + uint16_t rtt; +} __attribute__((packed)); + +#pragma pack(push, 1) +struct CONN_MGR_DIRECT_REQ { + uint8_t cmd; + uint8_t subcmd; + uint32_t request_id; + uint8_t addr_count; +}; + +struct CONN_MGR_DIRECT_RESP { + uint8_t cmd; + uint8_t subcmd; + uint32_t request_id; + uint8_t accepted; +}; + +struct CONN_MGR_INTERM_EXCHANGE_REQ { + uint8_t cmd; + uint8_t subcmd; + uint32_t request_id; + uint8_t candidate_count; + struct CONN_MGR_CANDIDATE candidates[4]; +}; + +struct CONN_MGR_INTERM_EXCHANGE_RESP { + uint8_t cmd; + uint8_t subcmd; + uint32_t request_id; + uint8_t my_count; + uint8_t your_count; + struct CONN_MGR_CANDIDATE my_candidates[4]; + struct CONN_MGR_CANDIDATE your_candidates[4]; +}; + +struct CONN_MGR_INTERM_SELECTED { + uint8_t cmd; + uint8_t subcmd; + uint32_t request_id; + uint8_t count; + struct CONN_MGR_CANDIDATE selected[3]; +}; + +struct CONN_MGR_DISCONNECT { + uint8_t cmd; + uint8_t subcmd; + uint64_t node_id; +}; +#pragma pack(pop) + +#define CONN_MGR_DIRECT_REQ_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_REQ) +#define CONN_MGR_DIRECT_RESP_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_RESP) +#define CONN_MGR_INTERM_EXCHANGE_REQ_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_REQ) +#define CONN_MGR_INTERM_EXCHANGE_RESP_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_RESP) +#define CONN_MGR_INTERM_SELECTED_SIZE sizeof(struct CONN_MGR_INTERM_SELECTED) +#define CONN_MGR_DISCONNECT_SIZE sizeof(struct CONN_MGR_DISCONNECT) + +typedef void (*conn_mgr_connect_callback_t)(int result, uint64_t node_id, void* arg); + +struct cm_cb_node { + conn_mgr_connect_callback_t cb; + void* arg; + struct cm_cb_node* next; +}; + +struct CONN_MGR_ENTRY { + uint64_t node_id; + uint8_t state; + uint8_t conn_type; + uint8_t alien; + uint32_t idle_timeout_ms; + uint64_t last_traffic_tb; + void* idle_timer; + + struct cm_cb_node* cb_list; + + uint64_t intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; + uint8_t intermediariy_count; + uint8_t rr_idx; + + uint8_t local_scan_state; + uint8_t main_connect_state; + struct { + uint8_t phase; + void* timer; + uint32_t request_id; + void* conn_ctx; + } main; + + struct CONN_MGR* mgr; +}; + +struct CONN_MGR { + struct UTUN_INSTANCE* instance; + struct CONN_MGR_ENTRY* entries; + size_t entry_count; + size_t entry_capacity; + + void* bg_ping_timer; + size_t bg_ping_cursor; + uint64_t bg_ping_cycle_start_tb; + + void* candidate_ping_timer; + + struct CONN_MGR_CANDIDATE best_candidates[CONN_MGR_MAX_CANDIDATES]; + uint8_t best_candidate_count; + + uint32_t next_request_id; + uint8_t initialized; + + struct cm_reverse_pending* reverse_pending; + struct cm_exchange_pending* exchange_pending; +}; + +struct CONN_MGR* conn_mgr_init(struct UTUN_INSTANCE* instance); +void conn_mgr_destroy(struct CONN_MGR* mgr); + +int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, + uint32_t idle_timeout_ms, + conn_mgr_connect_callback_t cb, void* cb_arg); +int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id); +int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t timeout_ms); +int conn_mgr_get_status(struct CONN_MGR* mgr, uint64_t node_id, + uint8_t* out_state, uint8_t* out_conn_type); +int conn_mgr_add_alien_node(struct CONN_MGR* mgr, const uint8_t* nodeinfo_data, size_t len); +int conn_mgr_send(struct CONN_MGR* mgr, uint64_t node_id, struct ll_entry* entry); + +void conn_mgr_update_best_candidates(struct CONN_MGR* mgr, uint64_t node_id, uint16_t rtt); + +#endif diff --git a/src/dummynet.c b/src/dummynet.c index 9a63a974..10c96049 100644 --- a/src/dummynet.c +++ b/src/dummynet.c @@ -577,6 +577,11 @@ struct dummynet_filter { uint32_t jitter_ms; struct ETCP_LINK* link; struct dummynet_stats stats; + struct { + uint32_t ip; + uint16_t port; + } block_addrs[4]; + uint8_t block_count; }; static ssize_t dummynet_filter_send_hook(socket_t fd, const void* buf, size_t len, @@ -585,6 +590,16 @@ static ssize_t dummynet_filter_send_hook(socket_t fd, const void* buf, size_t le struct dummynet_filter* df = (struct dummynet_filter*)ctx; (void)link; df->stats.recv++; + if (addr && addr->sa_family == AF_INET && df->block_count > 0) { + const struct sockaddr_in* sin = (const struct sockaddr_in*)addr; + for (uint8_t i = 0; i < df->block_count; i++) { + if (sin->sin_addr.s_addr == df->block_addrs[i].ip && + sin->sin_port == df->block_addrs[i].port) { + df->stats.lost++; + return (ssize_t)len; + } + } + } if (df->loss_permille > 0) { uint32_t roll = (uint32_t)rand() % 1000; if (roll < df->loss_permille) { @@ -637,6 +652,26 @@ void dummynet_filter_destroy(struct dummynet_filter* df) { u_free(df); } +void dummynet_filter_block_addr(struct dummynet_filter* df, uint32_t ip, uint16_t port) { + if (!df || df->block_count >= 4) return; + for (uint8_t i = 0; i < df->block_count; i++) + if (df->block_addrs[i].ip == ip && df->block_addrs[i].port == port) return; + df->block_addrs[df->block_count].ip = ip; + df->block_addrs[df->block_count].port = port; + df->block_count++; +} + +void dummynet_filter_unblock_addr(struct dummynet_filter* df, uint32_t ip, uint16_t port) { + if (!df) return; + for (uint8_t i = 0; i < df->block_count; i++) { + if (df->block_addrs[i].ip == ip && df->block_addrs[i].port == port) { + df->block_count--; + if (i < df->block_count) memmove(&df->block_addrs[i], &df->block_addrs[i+1], (df->block_count - i) * sizeof(df->block_addrs[0])); + return; + } + } +} + const struct dummynet_stats* dummynet_filter_get_stats(struct dummynet_filter* df) { return df ? &df->stats : NULL; } diff --git a/src/dummynet.h b/src/dummynet.h index 010b263b..dbc0bd8c 100644 --- a/src/dummynet.h +++ b/src/dummynet.h @@ -176,5 +176,7 @@ void dummynet_filter_attach(struct dummynet_filter* df, struct ETCP_LINK* link); void dummynet_filter_detach(struct dummynet_filter* df); void dummynet_filter_destroy(struct dummynet_filter* df); const struct dummynet_stats* dummynet_filter_get_stats(struct dummynet_filter* df); +void dummynet_filter_block_addr(struct dummynet_filter* df, uint32_t ip, uint16_t port); +void dummynet_filter_unblock_addr(struct dummynet_filter* df, uint32_t ip, uint16_t port); #endif // DUMMYNET_H diff --git a/src/etcp.h b/src/etcp.h index abc4de17..969f7a51 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -133,7 +133,7 @@ struct ETCP_CONN { // IDs and state - uint32_t next_tx_id; // Next TX ID + uint32_t next_tx_id; // ID для добавления в очередь отправки (с этим id будет добавлен следующий пакет) uint32_t last_rx_id; // Last received ID uint32_t last_delivered_id; // Last delivered to output_queue diff --git a/src/etcp_api.h b/src/etcp_api.h index 50017788..ee69c753 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -6,6 +6,8 @@ * - etcp_send() - отправить пакет в очередь normalizer * - etcp_bind() - подписаться на пакеты с определенным ID * - etcp_int_recv() - коллбэк для сбора пакетов из всех подключений + * !!! для обычной отправки-приёма между узлами используем более универсальный модушь etcp_router. + * это в первую очередь - апи для использования etcp-router-ом * * Формат кодограмм: * cmd = 0 - пакет для передачи адресату @@ -30,6 +32,7 @@ #define ETCP_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit) #define ETCP_ID_TCP_PROXY_CLIENT 0x07 // TCP proxy client (клиент) #define ETCP_ID_MSG_TRANSPORT 0x10 // msg_transport — локальный IPC транспорт сообщений +#define ETCP_ID_CONN_MGR 0x11 // Connection Manager — management connections // Forward declarations struct ETCP_CONN; diff --git a/src/etcp_router.c b/src/etcp_router.c index 9b46beae..a21b54f6 100644 --- a/src/etcp_router.c +++ b/src/etcp_router.c @@ -76,7 +76,18 @@ static int router_send_one_flags(struct ETCP_ROUTER_CONN* rconn, const uint8_t* } if (pl_len > 0) memcpy(dgram + SVC_ROUTE_HDR_SIZE, payload, pl_len); - struct ETCP_CONN* conn = route_bgp_find_conn_for_node(inst->bgp, rconn->remote_node_id); + struct ETCP_CONN* conn = NULL; + struct ROUTE_BGP* bgp = inst->bgp; + if (bgp) { + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, rconn->remote_node_id); + if (nq && nq->conn_mgr_type == CONN_TYPE_INDIRECT && nq->conn_mgr_intermediariy_count > 0) { + for (uint8_t i = 0; i < nq->conn_mgr_intermediariy_count; i++) { + conn = route_bgp_find_conn_for_node(bgp, nq->conn_mgr_intermediaries[i]); + if (conn) break; + } + } + } + if (!conn) conn = route_bgp_find_conn_for_node(bgp, rconn->remote_node_id); if (!conn) { DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router_send_one: no route to %016llx svc_id=%u", (unsigned long long)rconn->remote_node_id, rconn->svc_id); @@ -376,7 +387,7 @@ static void etcp_router_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) actual_payload = pl_len - SC_SIGN_SIZE; uint8_t* sig = pl + actual_payload; struct ROUTE_BGP* bgp = inst->bgp; - struct NODEINFO_Q* nq = bgp ? route_bgp_get_node(bgp, hdr->src_node_id) : NULL; + struct NODEINFO_Q* nq = bgp ? nodeinfo_find_by_id(bgp, hdr->src_node_id) : NULL; if (!nq) { DEBUG_WARN(DEBUG_CATEGORY_ETCPROUTE, "router: SIGNED packet from unknown node %016llx — dropping", (unsigned long long)hdr->src_node_id); queue_dgram_free(entry); queue_entry_free(entry); diff --git a/src/msg_transport.c b/src/msg_transport.c index 771a5c29..dcd009b3 100644 --- a/src/msg_transport.c +++ b/src/msg_transport.c @@ -462,7 +462,7 @@ static void handle_client_data(struct msg_transport* t, struct msg_client* clien send_error(client, MSG_ERR_PEER_NOT_FOUND, "BGP not initialized"); break; } - struct NODEINFO_Q* nq = route_bgp_get_node(bgp, peer_id); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, peer_id); if (!nq) { send_error(client, MSG_ERR_PEER_NOT_FOUND, "peer not found"); DEBUG_WARN(DEBUG_CATEGORY_MSGTRANSPORT, "GET_PEER_INFO: peer=%016llx not found", diff --git a/src/route_bgp.c b/src/route_bgp.c index f03f4633..beec94d5 100644 --- a/src/route_bgp.c +++ b/src/route_bgp.c @@ -370,9 +370,16 @@ void route_bgp_new_conn(struct ETCP_CONN* conn) { } struct ROUTE_BGP* bgp = conn->instance->bgp; + struct NODEINFO_Q* peer_nq = nodeinfo_find_by_id(bgp, conn->peer_node_id); route_bgp_add_to_senders(bgp, conn); + if (peer_nq && peer_nq->alien) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "route_bgp_new_conn: peer 0x%016llx is alien, skipping route exchange", + (unsigned long long)conn->peer_node_id); + return; + } + // Scan ALL connections and ALL links in instance - start NAT check if not started struct ETCP_CONN* c = conn->instance->connections; while (c) { @@ -483,19 +490,9 @@ void route_bgp_remove_conn(struct ETCP_CONN* conn) { * NEW NODEINFO BASED IMPLEMENTATION * ================================================ */ -struct NODEINFO_Q* route_bgp_get_node(struct ROUTE_BGP* bgp, uint64_t node_id) { - if (!bgp || !bgp->nodes) { - return NULL; - } - - uint64_t key = node_id; - struct ll_entry* e = queue_find_data_by_index(bgp->nodes, &key); - return e ? (struct NODEINFO_Q*)e : NULL; -} - struct ETCP_CONN* route_bgp_find_conn_for_node(struct ROUTE_BGP* bgp, uint64_t node_id) { if (!bgp) return NULL; - struct NODEINFO_Q* nq = route_bgp_get_node(bgp, node_id); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, node_id); if (!nq || !nq->paths || !nq->paths->head) return NULL; struct ll_entry* e = nq->paths->head; struct NODEINFO_PATH* best = NULL; @@ -744,7 +741,7 @@ int route_bgp_process_nodeinfo(struct ROUTE_BGP* bgp, struct ETCP_CONN* from, co return -1; } - struct NODEINFO_Q* nodeinfo1 = route_bgp_get_node(bgp, node_id); + struct NODEINFO_Q* nodeinfo1 = nodeinfo_find_by_id(bgp, node_id); uint8_t new_ver = ni->ver; if (nodeinfo1 && (int8_t)(nodeinfo1->last_ver-new_ver)>=0) { @@ -787,6 +784,8 @@ int route_bgp_process_nodeinfo(struct ROUTE_BGP* bgp, struct ETCP_CONN* from, co nodeinfo1->connectivity.interface_status = PROBE_RESULT_UNKNOWN; nodeinfo1->connectivity.nat_status = PROBE_RESULT_UNKNOWN; nodeinfo1->connectivity.real_status = PROBE_RESULT_UNKNOWN; + nodeinfo1->connectivity.ping_req_time = 0; + nodeinfo1->alien = 0; queue_data_put_with_index(bgp->nodes, &nodeinfo1->ll); } else { socks_changed = (nodeinfo1->node.local_v4_sockets != ni->local_v4_sockets) || @@ -872,7 +871,7 @@ int route_bgp_process_withdraw(struct ROUTE_BGP* bgp, struct ETCP_CONN* sender, uint64_t node_id = wp->node_id; uint64_t wd_source = wp->wd_source; - struct NODEINFO_Q* nq = route_bgp_get_node(bgp, node_id); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, node_id); if (!nq) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "node not found"); return 0; diff --git a/src/route_bgp.h b/src/route_bgp.h index ec8d9e0c..4685d45f 100644 --- a/src/route_bgp.h +++ b/src/route_bgp.h @@ -155,13 +155,6 @@ void route_bgp_send_nodeinfo(struct NODEINFO_Q* node, struct ETCP_CONN* conn); */ void route_bgp_send_withdraw(struct ROUTE_BGP* bgp, uint64_t node_id); -/** - * @brief Поиск NODEINFO_Q по node_id через hash в nodes queue. - * - * @return node или NULL - */ -struct NODEINFO_Q* route_bgp_get_node(struct ROUTE_BGP* bgp, uint64_t node_id); - int nodeinfo_dyn_size(struct NODEINFO* node); /** diff --git a/src/route_node.c b/src/route_node.c index f43420c2..b7d23ed0 100644 --- a/src/route_node.c +++ b/src/route_node.c @@ -121,6 +121,13 @@ int get_node_v6_routes(struct NODEINFO_Q *node, const struct NODEINFO_IPV6_SUBNE return (int)info->local_v6_subnets; } +struct NODEINFO_Q* nodeinfo_find_by_id(struct ROUTE_BGP* bgp, uint64_t node_id) { + if (!bgp || !bgp->nodes) return NULL; + uint64_t key = node_id; + struct ll_entry* e = queue_find_data_by_index(bgp->nodes, &key); + return e ? (struct NODEINFO_Q*)e : NULL; +} + int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BGP* bgp) { if (!instance || !bgp) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_bgp_update_my_nodeinfo: invalid args"); @@ -181,7 +188,7 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG } } - if (changed) { + if (changed) { if (bgp->local_node) u_free(bgp->local_node); bgp->local_node = u_calloc(1, sizeof(struct NODEINFO_Q) + dyn); if (!bgp->local_node) return -1; diff --git a/src/route_node.h b/src/route_node.h index 0ec2b52e..4159d96a 100644 --- a/src/route_node.h +++ b/src/route_node.h @@ -33,6 +33,12 @@ struct UTUN_INSTANCE; #define CONN_PROBE_COUNT 3 #define CONN_PROBE_TIMEOUT_MS 1000 +#define CONN_TYPE_NONE 0 +#define CONN_TYPE_DIRECT 1 +#define CONN_TYPE_REVERSE 2 +#define CONN_TYPE_INDIRECT 3 +#define CONN_MGR_MAX_INTERMEDIARIES 3 + // ---- состояние связности с удалённым узлом (локальное, не передаётся по BGP) ---- struct NODE_CONNECTIVITY { @@ -49,6 +55,7 @@ struct NODE_CONNECTIVITY { uint64_t interface_probe_time; uint64_t nat_probe_time; uint64_t real_probe_time; + uint64_t ping_req_time; // время последнего полученного ping request от этого узла (0.1ms tb) }; /** @@ -133,6 +140,10 @@ struct NODEINFO_Q { struct ll_queue* paths; // сюда помещаем struct NODEINFO_PATH uint8_t dirty; uint8_t last_ver; + uint8_t alien; // 1 = чужой узел (не запрашиваем роутинг при connect) + uint8_t conn_mgr_type; // CONN_TYPE_* — тип managed-подключения + uint64_t conn_mgr_intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; // посредники для INDIRECT + uint8_t conn_mgr_intermediariy_count; struct NODE_CONNECTIVITY connectivity; // состояние связности (локальное) struct ETCP_SOCKET* best_socket; struct NODEINFO node; // Всегда в конце структуры - динамически расширяемый блок @@ -174,4 +185,6 @@ int get_node_v6_sockets_meta(struct NODEINFO_Q *node, const struct NODEINFO_IPV6 */ int get_node_v6_addrs(struct NODEINFO_Q *node, const struct NODEINFO_IPV6_ADDR **out_addrs); +struct NODEINFO_Q* nodeinfo_find_by_id(struct ROUTE_BGP* bgp, uint64_t node_id); + #endif // ROUTE_NODE_H diff --git a/src/route_ping.c b/src/route_ping.c index 0a276226..f9a7acdc 100644 --- a/src/route_ping.c +++ b/src/route_ping.c @@ -325,6 +325,11 @@ void route_ping_handle_req(struct ROUTE_BGP* bgp, ctx->request_id = req_pkt->request_id; ctx->count_total = req_pkt->count; ctx->timeout_ms = req_pkt->timeout_ms; + + { + struct NODEINFO_Q* req_nq = nodeinfo_find_by_id(bgp, from_conn->peer_node_id); + if (req_nq) req_nq->connectivity.ping_req_time = get_time_tb(); + } if (len >= sizeof(struct BGP_PING_REQUEST)) { memcpy(ctx->pubkey, req_pkt->pubkey, SC_PUBKEY_SIZE); } @@ -338,7 +343,7 @@ void route_ping_handle_req(struct ROUTE_BGP* bgp, /* Если target не указан — можно разрешить из nodeinfo */ /* if (sin->sin_addr.s_addr == 0 && sin->sin_port == 0) { - struct NODEINFO_Q* nq = route_bgp_get_node(bgp, req_pkt->node_id); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(bgp, req_pkt->node_id); if (nq) { const struct NODEINFO_IPV4_ADDR* addrs; int count = get_node_v4_addrs(nq, &addrs); diff --git a/src/utun_instance.c b/src/utun_instance.c index 383cf651..114b267b 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -11,6 +11,7 @@ #include "route_bgp.h" #include "etcp_connections.h" #include "etcp.h" +#include "conn_mgr.h" #include "control_server.h" #include "msg_transport.h" #include "../lib/u_async.h" @@ -129,6 +130,9 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u if (route_insert(instance->rt,instance->bgp->local_node)) DEBUG_INFO(DEBUG_CATEGORY_ROUTING,"Added local routes"); else DEBUG_WARN(DEBUG_CATEGORY_ROUTING,"Failed to add local routes"); } } + + instance->conn_mgr = conn_mgr_init(instance); + if (!instance->conn_mgr) DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "Failed to initialize Connection Manager (non-fatal)"); // Initialize firewall fw_init(&instance->fw); @@ -373,6 +377,8 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Cleanup etcp_router etcp_router_destroy(instance); + + if (instance->conn_mgr) { conn_mgr_destroy(instance->conn_mgr); instance->conn_mgr = NULL; } // Cleanup BGP module if (instance->bgp) { diff --git a/src/utun_instance.h b/src/utun_instance.h index 5872f811..6d663009 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -33,6 +33,7 @@ struct ROUTE_BGP; struct control_server; struct msg_transport; struct PING_CONTEXT; +struct CONN_MGR; // uTun instance configuration struct UTUN_INSTANCE { @@ -112,6 +113,7 @@ struct UTUN_INSTANCE { // etcp_router bindings и seq-connections (per-instance service routing) struct ETCP_ROUTER_BINDINGS router_bindings; struct ll_queue* router_conns; + struct CONN_MGR* conn_mgr; // Connection Manager (может быть NULL) // TCP proxy server (exit node) struct tcp_proxy_server tcp_proxy_server; diff --git a/tests/Makefile.am b/tests/Makefile.am index e3e34246..d26a7a61 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -44,6 +44,7 @@ check_PROGRAMS = \ test_icmp_proxy \ test_tcp_proxy_client \ test_bgp_route_exchange \ + test_conn_mgr \ test_bbr_integration \ test_intensive_memory_pool \ test_tcp_io \ @@ -94,6 +95,7 @@ ETCP_FULL_OBJS = \ $(top_builddir)/src/utun-route_node.o \ $(top_builddir)/src/utun-route_node_lmdb.o \ $(top_builddir)/src/utun-route_connectivity.o \ + $(top_builddir)/src/utun-conn_mgr.o \ $(top_builddir)/src/utun-routing.o \ $(top_builddir)/src/utun-tun_if.o \ $(top_builddir)/src/utun-tun_route.o \ @@ -308,6 +310,10 @@ test_bgp_route_exchange_SOURCES = test_bgp_route_exchange.c test_bgp_route_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_bgp_route_exchange_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_conn_mgr_SOURCES = test_conn_mgr.c +test_conn_mgr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_conn_mgr_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_bbr_integration_SOURCES = bbr_integration/test_bbr_integration.c test_bbr_integration_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_bbr_integration_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_conn_mgr.c b/tests/test_conn_mgr.c new file mode 100644 index 00000000..1a8ce2c8 --- /dev/null +++ b/tests/test_conn_mgr.c @@ -0,0 +1,136 @@ +/** + * @file test_conn_mgr.c + * @brief Connection Manager — smoke test + * + * Test 1: Direct via existing ETCP conn → CONN_TYPE_DIRECT + * Test 2: Status check → state/type + * Test 3: Alien node → NODEINFO_Q.alien=1 + * + * Pending: + * Test 4: New link UP callback (INIT to new peer via conn_mgr) + * Test 5: Indirect through intermediary (needs full Phase 3 exchange) + * Test 6: Idle timeout + */ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifndef _WIN32 +#include +#endif + +#include "../src/etcp.h" +#include "../src/etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "../src/route_bgp.h" +#include "../src/conn_mgr.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TIMEOUT_TB 300000 +#define POLL_MS 5 +#define NID_A 0xAAAA000000000001ULL +#define NID_B 0xBBBB000000000002ULL + +static struct UTUN_INSTANCE* g_a = NULL, *g_b = NULL; +static struct UASYNC* ua = NULL; +static int result = 0; +static void* ttimer = NULL; +static char tdir[] = "/tmp/utun_cm_XXXXXX"; +static char ca[256], cb[256]; +static int pa = 0, pb = 0; +static volatile int cdone = 0, cresult = 0; + +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 char* gv(const char* p, const char* k) { + struct utun_config* c = parse_config(p); if (!c) return NULL; + char* r = (strcmp(k, "pub") == 0) ? u_strdup(c->global.my_public_key_hex) : u_strdup(c->global.my_private_key_hex); + free_config(c); return r; +} +static int lks(struct UTUN_INSTANCE* i) { + int n = 0; struct ETCP_CONN* c = i->connections; + while (c) { struct ETCP_LINK* l = c->links; while (l) { if (l->initialized && l->link_status) n++; l = l->next; } c = c->next; } + return n; +} +static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 2; } +static void test1(void* arg); +static void test2(void* arg); +static void test3(void* arg); + +static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); result = 2; } +static void ccb(int r, uint64_t id, void* arg) { + (void)arg; fprintf(stderr, "connect_cb: result=%d node=0x%llx\n", r, (unsigned long long)id); fflush(stderr); + cdone = 1; cresult = r; +} + +static void test1(void* arg) { + (void)arg; if (result) return; + if (lks(g_a) < 1 || lks(g_b) < 1) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1"); return; } + fprintf(stderr, "Test 1: direct — conn_mgr_connect_node(B)\n"); fflush(stderr); + conn_mgr_connect_node(g_a->conn_mgr, NID_B, 0, ccb, NULL); + uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2, "t2a"); +} +static void test2(void* arg) { + (void)arg; if (result) return; + if (!cdone) { uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2, "t2b"); return; } + if (cresult != CONN_MGR_OK) { fail("connect failed"); return; } + uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, NID_B, &st, &ty); + fprintf(stderr, "Test 2: status state=%d type=%d\n", st, ty); fflush(stderr); + if (st == CONN_MGR_STATE_CONNECTED && ty == CONN_TYPE_DIRECT) fprintf(stderr, " OK: DIRECT\n"); + else { fail("status mismatch"); return; } + test3(NULL); +} +static void test3(void* arg) { + (void)arg; + uint64_t aid = 0xEEEE000000000001ULL; + struct NODEINFO ni; memset(&ni, 0, sizeof(ni)); ni.node_id = aid; ni.ver = 1; + conn_mgr_add_alien_node(g_a->conn_mgr, (uint8_t*)&ni, sizeof(ni)); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(g_a->bgp, aid); + fprintf(stderr, "Test 3: alien alien=%d\n", nq ? nq->alien : -1); fflush(stderr); + if (nq && nq->alien == 1) fprintf(stderr, " OK: alien flag set\n"); + else { fail("alien flag not set"); return; } + fprintf(stderr, "=== ALL DONE ===\n"); fflush(stderr); + result = (result == 0) ? 1 : 2; +} + +static void setup(void) { + test_mkdtemp(tdir); + int base = 47000 + (getpid() % 15000); pa = base; pb = base + 1; + snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir); + wf(ca, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_A, pa); + wf(cb, "[global]\nmy_node_id=0x%llx\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, pb); + config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); + char *p0 = gv(ca,"pub"), *r0 = gv(ca,"priv"), *p1 = gv(cb,"pub"), *r1 = gv(cb,"priv"); + wf(ca, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.1/24\ntun_ifname=tun99\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", NID_A, r0, p0, pa, p1, pb); + wf(cb, "[global]\nmy_node_id=0x%llx\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", NID_B, r1, p1, pb); + u_free(p0); u_free(r0); u_free(p1); u_free(r1); +} +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); setup(); + ua = uasync_create(); + g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb); + if (!g_a || !g_b) goto done; + utun_instance_init(g_a); utun_instance_init(g_b); + uasync_call_soon(ua, NULL, test1); + ttimer = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to"); + { int el = 0; while (!result && el < TIMEOUT_TB + 5000) { uasync_poll(ua, POLL_MS); el += POLL_MS; } } + fprintf(stderr, "final result=%d\n", result); fflush(stderr); +done: + if (ttimer) uasync_cancel_timeout(ua, ttimer); + 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); ua = NULL; } + cleanup(); + return (result == 1) ? 0 : 1; +} diff --git a/tests/test_ipv6_sockets.c b/tests/test_ipv6_sockets.c index 3d3644b8..30eac267 100644 --- a/tests/test_ipv6_sockets.c +++ b/tests/test_ipv6_sockets.c @@ -221,7 +221,7 @@ static int verify_remote_v6_nodeinfo(const char* name, struct UTUN_INSTANCE* ins printf("FAIL [%s]: no bgp\n", name); return 0; } - struct NODEINFO_Q* nq = route_bgp_get_node(inst->bgp, peer_node_id); + struct NODEINFO_Q* nq = nodeinfo_find_by_id(inst->bgp, peer_node_id); if (!nq) { printf("FAIL [%s]: remote node %016llx not found\n", name, (unsigned long long)peer_node_id); return 0; diff --git a/tests/test_nat_detection.c b/tests/test_nat_detection.c index 05a16d23..5dfe8ba3 100644 --- a/tests/test_nat_detection.c +++ b/tests/test_nat_detection.c @@ -244,16 +244,16 @@ int main(void) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "Waiting for BGP exchange (S learns C1 and C2)..."); int bgp_wait_cycles = 0; while (!test_timed_out && bgp_wait_cycles < 500) { - if (inst_s->bgp && route_bgp_get_node(inst_s->bgp, NODE_ID_C1) != NULL && - route_bgp_get_node(inst_s->bgp, NODE_ID_C2) != NULL) { + if (inst_s->bgp && nodeinfo_find_by_id(inst_s->bgp, NODE_ID_C1) != NULL && + nodeinfo_find_by_id(inst_s->bgp, NODE_ID_C2) != NULL) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "S learned about C1 and C2"); break; } uasync_poll(ua, 10); bgp_wait_cycles++; } - if (!inst_s->bgp || route_bgp_get_node(inst_s->bgp, NODE_ID_C1) == NULL || - route_bgp_get_node(inst_s->bgp, NODE_ID_C2) == NULL) { + if (!inst_s->bgp || nodeinfo_find_by_id(inst_s->bgp, NODE_ID_C1) == NULL || + nodeinfo_find_by_id(inst_s->bgp, NODE_ID_C2) == NULL) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "S did not learn about C1/C2 in time"); goto cleanup; } @@ -313,7 +313,7 @@ int main(void) { } // 5. Verify all NAT fields on server - struct NODEINFO_Q* node_c1 = route_bgp_get_node(inst_s->bgp, NODE_ID_C1); + struct NODEINFO_Q* node_c1 = nodeinfo_find_by_id(inst_s->bgp, NODE_ID_C1); if (!node_c1) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "NODEINFO_Q for C1 disappeared"); goto cleanup; @@ -408,7 +408,7 @@ int main(void) { bgp_wait_cycles = 0; int c2_verified_nat = 0; while (!test_timed_out && bgp_wait_cycles < 500) { - struct NODEINFO_Q* node_c1_on_c2 = inst_c2->bgp ? route_bgp_get_node(inst_c2->bgp, NODE_ID_C1) : NULL; + struct NODEINFO_Q* node_c1_on_c2 = inst_c2->bgp ? nodeinfo_find_by_id(inst_c2->bgp, NODE_ID_C1) : NULL; if (node_c1_on_c2) { const struct NODEINFO_IPV4_SOCKET_META* meta = NULL; int meta_count = get_node_v4_sockets_meta(node_c1_on_c2, &meta); @@ -477,14 +477,14 @@ int main(void) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "Waiting for C2 to learn C1 nodeinfo..."); bgp_wait_cycles = 0; while (!test_timed_out && bgp_wait_cycles < 500) { - if (inst_c2->bgp && route_bgp_get_node(inst_c2->bgp, NODE_ID_C1) != NULL) { + if (inst_c2->bgp && nodeinfo_find_by_id(inst_c2->bgp, NODE_ID_C1) != NULL) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "C2 learned C1 nodeinfo"); break; } uasync_poll(ua, 10); bgp_wait_cycles++; } - if (!inst_c2->bgp || route_bgp_get_node(inst_c2->bgp, NODE_ID_C1) == NULL) { + if (!inst_c2->bgp || nodeinfo_find_by_id(inst_c2->bgp, NODE_ID_C1) == NULL) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "FAIL: C2 did not learn C1 nodeinfo in time"); goto cleanup; } diff --git a/tests/test_nat_transport.c b/tests/test_nat_transport.c index c18fc446..62d582a3 100644 --- a/tests/test_nat_transport.c +++ b/tests/test_nat_transport.c @@ -565,8 +565,8 @@ int main(void) { DEBUG_INFO(DEBUG_CATEGORY_NAT, "Waiting for BGP..."); int bgp_cycles = 0; while (!test_timed_out && bgp_cycles < 500) { - if (inst_provider->bgp && route_bgp_get_node(inst_provider->bgp, NODE_ID_CLIENT) && - inst_client->bgp && route_bgp_get_node(inst_client->bgp, NODE_ID_PROVIDER)) { + if (inst_provider->bgp && nodeinfo_find_by_id(inst_provider->bgp, NODE_ID_CLIENT) && + inst_client->bgp && nodeinfo_find_by_id(inst_client->bgp, NODE_ID_PROVIDER)) { DEBUG_INFO(DEBUG_CATEGORY_NAT, "BGP exchanged"); break; } diff --git a/tests/test_route_ping.c b/tests/test_route_ping.c index 193101ed..be34c968 100644 --- a/tests/test_route_ping.c +++ b/tests/test_route_ping.c @@ -248,14 +248,14 @@ int main(void) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "Waiting for BGP exchange (B learns C)..."); int bgp_wait_cycles = 0; while (!test_timed_out && bgp_wait_cycles < 500) { - if (inst_b->bgp && route_bgp_get_node(inst_b->bgp, NODE_ID_C) != NULL) { + if (inst_b->bgp && nodeinfo_find_by_id(inst_b->bgp, NODE_ID_C) != NULL) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "B learned about C"); break; } uasync_poll(ua, 10); bgp_wait_cycles++; } - if (!inst_b->bgp || route_bgp_get_node(inst_b->bgp, NODE_ID_C) == NULL) { + if (!inst_b->bgp || nodeinfo_find_by_id(inst_b->bgp, NODE_ID_C) == NULL) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "B did not learn about C in time"); goto cleanup; } @@ -271,7 +271,7 @@ int main(void) { /* addr[] в network byte order, port в host byte order */ uint32_t target_ip = 0; uint16_t target_port = 0; - struct NODEINFO_Q* nq = inst_b->bgp ? route_bgp_get_node(inst_b->bgp, NODE_ID_C) : NULL; + struct NODEINFO_Q* nq = inst_b->bgp ? nodeinfo_find_by_id(inst_b->bgp, NODE_ID_C) : NULL; if (nq) { const struct NODEINFO_IPV4_ADDR* addrs; int sc = get_node_v4_addrs(nq, &addrs);