You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

760 lines
45 KiB

/* conn_mgr_core.c — init/destroy, open/close/handles, DIRECT, REVERSE, invite */
#include "conn_mgr_priv.h"
#include <stdlib.h>
#include <string.h>
#ifdef _WIN32
#include <winsock2.h>
#include <ws2tcpip.h>
#else
#include <arpa/inet.h>
#endif
#include "../lib/platform_compat.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "utun_instance.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "etcp_router.h"
#include "config_parser.h"
#include "topo_node.h"
#include "topo_node_sqlite.h"
#include "topo_group.h"
#include "route_ping.h"
#ifdef UTUN_HAVE_STANDBY
#include "standby.h"
#endif
/* ═══════ forward-декларации ═══════ */
void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry);
/* ═══════ утилиты ═══════ */
static struct CONN_MGR_HANDLE* cm_handle_new(struct CONN_MGR_ENTRY* entry, uint64_t group_id,
conn_mgr_cb_t cb, void* cb_arg) {
struct CONN_MGR_HANDLE* h = u_calloc(1, sizeof(*h));
if (!h) return NULL;
h->node_id = entry->node_id; h->group_id = group_id; h->cb = cb; h->cb_arg = cb_arg;
h->entry = entry; h->next = entry->handles; entry->handles = h;
return h;
}
static void cm_deliver_up_cb(void* arg) {
struct CONN_MGR_HANDLE* h = (struct CONN_MGR_HANDLE*)arg;
h->deliver_up_token = NULL;
if (h->cb) h->cb(h, h->node_id, h->group_id, CONN_EVENT_UP, h->cb_arg);
}
/* Рассылает событие ВСЕМ handle'ам entry. Одно соединение — много слушателей:
* каждый conn_mgr_open создал свой handle, все получают UP/DOWN/TIMEOUT.
* Все поля handle'а читаются до вызова коллбэка — коллбэк может закрыть
* handle (conn_mgr_close), после вызова к handle не обращаемся. */
void cm_deliver_event(struct CONN_MGR_ENTRY* entry, enum conn_mgr_event event) {
if (!entry || !entry->mgr || !entry->mgr->group) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "cm_deliver_event: BAD entry=%p mgr=%p", entry, entry ? (void*)entry->mgr : NULL);
return;
}
uint64_t gid = entry->mgr->group->group_id;
struct CONN_MGR_HANDLE* h = entry->handles;
while (h) {
struct CONN_MGR_HANDLE* next = h->next;
uint64_t node_id = h->node_id;
conn_mgr_cb_t cb = h->cb;
void* cb_arg = h->cb_arg;
if (cb) cb(h, node_id, gid, event, cb_arg);
h = next;
}
}
/* ═══════ NAT / адресные утилиты ═══════ */
/* Проверяет NAT-совместимость нашего сокета с целевым: возвращает 1 (публичный/EIM —
* можно пробить), 2 (локальная сеть — оба PRIVATE), 0 (несовместим). Используется
* в DIRECT фазе чтобы не создавать заведомо непробиваемые линки. */
int cm_nat_compatible(struct ETCP_SOCKET* our, uint8_t tcfg, uint8_t tnat) {
uint8_t tt = tnat >= NAT_VERIFIED_UNKNOWN ? CFG_SERVER_TYPE_PUBLIC : tcfg;
uint8_t ot = our->type;
if (ot == CFG_SERVER_TYPE_PUBLIC || ot == CFG_SERVER_TYPE_UNKNOWN) return 1;
if (tt == CFG_SERVER_TYPE_PUBLIC) return 1; /* цель публичная/проверенная — исходящее работает из-за любого NAT */
if (ot == CFG_SERVER_TYPE_NAT && our->nat_type == NAT_VERIFIED_EIM && tt == CFG_SERVER_TYPE_NAT) return 1;
if ((ot == CFG_SERVER_TYPE_PRIVATE || ot == CFG_SERVER_TYPE_LOCAL) && tt == CFG_SERVER_TYPE_PRIVATE) return 2;
return 0;
}
/* Классификация IPv6-адреса: link-local (fe80), локальный ULA (fc/fd),
* глобальный/прямой, другое (loopback, multicast, ::). */
uint8_t cm_classify_v6_addr(const uint8_t addr[16]) {
if (addr[0] == 0xfe && (addr[1] & 0xc0) == 0x80) return CM_V6_LL;
if (addr[0] == 0xfc || addr[0] == 0xfd) return CM_V6_LOC;
if (addr[0] == 0xff) return CM_V6_OTH;
{ static const uint8_t z[16]; if (!memcmp(addr, z, 16)) return CM_V6_OTH; }
{ static const 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)) return CM_V6_OTH; }
return CM_V6_DIR;
}
uint8_t cm_sock_v6_classify(const struct ETCP_SOCKET* s) {
if (s->local_addr.ss_family != AF_INET6) return CM_V6_OTH;
const uint8_t* a6 = ((const struct sockaddr_in6*)&s->local_addr)->sin6_addr.s6_addr;
static const uint8_t z[16];
if (!memcmp(a6, z, 16)) return CM_V6_ANY;
return cm_classify_v6_addr(a6);
}
void cm_add_v4_link(struct ETCP_CONN* conn, const uint8_t* addr, uint16_t port, struct ETCP_SOCKET* s) {
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, addr, 4); sin.sin_port = port;
struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
etcp_link_new(conn, s, &sa, 0);
}
void cm_add_v6_link(struct ETCP_CONN* conn, const uint8_t addr[16], uint16_t port, struct ETCP_SOCKET* s) {
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, addr, 16); sin6.sin6_port = htons(port);
if (cm_classify_v6_addr(addr) == CM_V6_LL) sin6.sin6_scope_id = s->netif_index;
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
etcp_link_new(conn, s, &sa, 0);
}
/* Проверяет есть ли у узла прямые IP (public/EIM/DIRECT на NAT/INTERFACE адресах) —
* можем ли мы достучаться до узла напрямую. Определяет, пойдём в REVERSE или INDIRECT. */
int cm_has_direct_ip(struct CONN_MGR* mgr, struct TOPO_GROUP_NODE* nq) {
if (!nq) return 0;
struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, nq->node_id);
if (!ni) return 0;
const uint8_t pub_conf[2] = {CFG_SERVER_TYPE_PUBLIC, CFG_SERVER_TYPE_UNKNOWN};
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = ni->v4_sock_meta; m; m = m->next)
if (m->id == a->socket_id && (m->config_type == pub_conf[0] || m->config_type == pub_conf[1]
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT)) return 1;
}
for (const struct TOPO_ADDR6* a6 = ni->v6_addrs; a6; a6 = a6->next) {
if ((a6->type != TOPO_ADDR_NAT && a6->type != TOPO_ADDR_INTERFACE) || cm_classify_v6_addr(a6->addr) != CM_V6_DIR) continue;
for (const struct TOPO_SOCKMETA6* m = ni->v6_sock_meta; m; m = m->next)
if (m->id == a6->socket_id && (m->config_type == pub_conf[0] || m->config_type == pub_conf[1]
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT)) return 1;
}
return 0;
}
int cm_has_local_addr(struct CONN_MGR* mgr, struct TOPO_GROUP_NODE* nq) {
if (!nq) return 0;
struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, nq->node_id);
if (!ni) return 0;
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = ni->v4_sock_meta; m; m = m->next)
if (m->id == a->socket_id && (m->config_type == CFG_SERVER_TYPE_PRIVATE || m->config_type == CFG_SERVER_TYPE_LOCAL)) return 1;
}
for (const struct TOPO_ADDR6* a6 = ni->v6_addrs; a6; a6 = a6->next) {
if (a6->type != TOPO_ADDR_INTERFACE) continue;
uint8_t cls = cm_classify_v6_addr(a6->addr);
if (cls != CM_V6_LL && cls != CM_V6_LOC) continue;
for (const struct TOPO_SOCKMETA6* m = ni->v6_sock_meta; m; m = m->next)
if (m->id == a6->socket_id && (m->config_type == CFG_SERVER_TYPE_PRIVATE || m->config_type == CFG_SERVER_TYPE_LOCAL)) return 1;
}
return 0;
}
/* ═══════ entries ═══════ */
struct CONN_MGR_ENTRY* cm_find_entry(struct CONN_MGR* mgr, uint64_t node_id) {
if (!mgr || !mgr->entries) return NULL;
struct ll_entry* e = queue_find_data_by_index(mgr->entries, &node_id);
return e ? (struct CONN_MGR_ENTRY*)e : NULL;
}
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;
struct ll_entry* qe = queue_entry_new(sizeof(struct CONN_MGR_ENTRY) - sizeof(struct ll_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: queue_entry_new failed"); return NULL; }
struct CONN_MGR_ENTRY* new_e = (struct CONN_MGR_ENTRY*)qe;
memset((uint8_t*)new_e + sizeof(struct ll_entry), 0, sizeof(*new_e) - sizeof(struct ll_entry));
new_e->node_id = node_id; new_e->mgr = mgr;
queue_data_put_with_index(mgr->entries, qe);
return new_e;
}
void cm_entry_cleanup(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; }
{ struct cm_exchange_pending** pp = &entry->mgr->exchange_pending;
while (*pp) {
if ((*pp)->entry == entry) {
if ((*pp)->timer && (*pp)->timer != entry->main.timer)
{ uasync_cancel_timeout(entry->mgr->instance->ua, (*pp)->timer); (*pp)->timer = NULL; }
struct cm_exchange_pending* ep = *pp; *pp = ep->next; u_free(ep);
} else pp = &(*pp)->next;
}
}
if (entry->ncd_handle) { node_conn_direct_force_close(entry->ncd_handle); entry->ncd_handle = NULL; }
entry->handles = NULL;
cm_clear_nodeinfo(entry->mgr, entry->node_id);
entry->state = CONN_MGR_STATE_DISCONNECTED; entry->conn_type = CONN_TYPE_NONE;
{ struct CONN_MGR* m = entry->mgr;
queue_remove_data(m->entries, &entry->ll);
queue_entry_free(&entry->ll); }
}
void cm_clear_nodeinfo(struct CONN_MGR* mgr, uint64_t node_id) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(mgr->group, node_id);
if (!nq) return;
nq->conn_mgr_type = CONN_TYPE_NONE; nq->conn_mgr_intermediariy_count = 0;
}
void cm_update_nodeinfo(struct CONN_MGR_ENTRY* entry) {
if (!entry->mgr->group) return;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(entry->mgr->group, 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;
if (entry->conn_type == CONN_TYPE_INDIRECT) { nq->conn_presence |= NCONN_INDIRECT; nq->conn_up |= NCONN_INDIRECT; }
else { nq->conn_presence &= ~NCONN_INDIRECT; nq->conn_up &= ~NCONN_INDIRECT; }
topo_fire_nodeinfo_cbk(entry->mgr->instance, entry->mgr->group, nq);
}
void cm_cleanup_db_node(struct CONN_MGR_ENTRY* entry) {
if (!entry || !entry->mgr || !entry->mgr->group || !entry->db_loaded) return;
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, entry->node_id);
if (!nq) return;
if (nq->paths && queue_entry_count(nq->paths) > 0) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: db_node 0x%016llx got paths, keep in group", (unsigned long long)entry->node_id);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: removing db_node 0x%016llx from group (no paths)", (unsigned long long)entry->node_id);
queue_remove_data(group->nodes, &nq->ll);
topo_nodeq_free_group_fields(entry->mgr->instance->topo_groups, nq);
queue_entry_free(&nq->ll);
}
/* ═══════ NCD коллбэки ═══════ */
/* NCD (node conn direct) события → conn_mgr для DIRECT/REVERSE фаз:
* UP — соединение поднялось, помечаем CONNECTED, доставляем UP всем handle'ам.
* TIMEOUT — NCD таймаут истёк, переходим к REVERSE или INDIRECT, либо фейлим db_node.
* DOWN — соединение упало, доставляем DOWN. */
void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; (void)ncd_h;
if (!entry || !entry->mgr || entry->state == CONN_MGR_STATE_DISCONNECTED) return;
switch (ncd_ev) {
case NCD_EVENT_UP:
if (entry->state == CONN_MGR_STATE_CONNECTED) return;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP handshake OK with 0x%016llx — connection ESTABLISHED",
(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_event(entry, CONN_EVENT_UP);
break;
case NCD_EVENT_TIMEOUT:
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP connect TIMEOUT to 0x%016llx — entry=%p mgr=%p db_loaded=%d",
(unsigned long long)entry->node_id, entry, (void*)entry->mgr, entry->db_loaded);
if (entry->db_loaded) {
entry->main_connect_state = CM_TRY_FAILED;
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: TIMEOUT after force_close, calling cleanup");
cm_cleanup_db_node(entry);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: TIMEOUT after cleanup, calling deliver");
cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: TIMEOUT done");
return;
}
entry->main_connect_state = CM_TRY_FAILED;
{ struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
if (entry->local_scan_state != CM_TRY_OK) {
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; }
cm_start_phase_indirect(entry); return;
}
}
if (entry->local_scan_state == CM_TRY_FAILED) cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
break;
case NCD_EVENT_DOWN:
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "link: ETCP link DROPPED to 0x%016llx",
(unsigned long long)entry->node_id);
cm_deliver_event(entry, CONN_EVENT_DOWN);
break;
}
}
/* ═══════ init / destroy ═══════ */
struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) {
if (!group || !group->instance) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NULL group"); return NULL; }
struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR));
if (!mgr) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "alloc failed"); return NULL; }
mgr->instance = group->instance; mgr->group = group;
mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS;
mgr->entries = queue_new(mgr->instance->ua, 256,
offsetof(struct CONN_MGR_ENTRY, node_id) - sizeof(struct ll_entry),
8, "conn_mgr_entries");
if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "entries alloc failed"); u_free(mgr); return NULL; }
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");
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");
mgr->bg_ping_cycle_start_tb = get_time_tb(); mgr->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: initialized for group %016llx", (unsigned long long)group->group_id);
return mgr;
}
void conn_mgr_destroy(struct CONN_MGR* mgr) {
if (!mgr) return;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 0 enter mgr=%p entries=%p exch=%p rev=%p",
mgr, mgr->entries, mgr->exchange_pending, mgr->reverse_pending);
if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; }
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 1 bg_ping done");
if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; }
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 2 cand_ping done");
#ifdef UTUN_HAVE_STANDBY
if (mgr->bg_ping_wait) { standby_wait_cancel(mgr->bg_ping_wait); mgr->bg_ping_wait = NULL; }
if (mgr->candidate_ping_wait) { standby_wait_cancel(mgr->candidate_ping_wait); mgr->candidate_ping_wait = NULL; }
#endif
{ size_t ec = queue_entry_count(mgr->entries);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 3 entries count=%zu", ec);
struct ll_entry* e = mgr->entries->head;
int clean = 0;
while (e) { struct ll_entry* next = e->next; struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)e; cm_entry_cleanup(entry); clean++; e = next; }
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 3 entries cleaned (%d)", clean);
while (mgr->exchange_pending) {
struct cm_exchange_pending* ep = mgr->exchange_pending; mgr->exchange_pending = ep->next;
if (ep->timer) uasync_cancel_timeout(mgr->instance->ua, ep->timer);
u_free(ep);
}
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 5 exchange_pending done");
while (mgr->reverse_pending) { struct cm_reverse_pending* rp = mgr->reverse_pending; mgr->reverse_pending = rp->next; u_free(rp); }
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 6 reverse_pending done");
queue_free(mgr->entries); mgr->initialized = 0;
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[CM_DESTROY] 7 queue_free done (was %zu)", ec); }
u_free(mgr);
}
/* ═══════ публичное API ═══════ */
int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms,
conn_mgr_cb_t cb, void* cb_arg, struct CONN_MGR_HANDLE** out_handle) {
if (out_handle) *out_handle = NULL;
if (!mgr) { if (cb) cb(NULL, node_id, 0, CONN_EVENT_TIMEOUT, cb_arg); return -1; }
struct TOPO_GROUP* group = mgr->group;
if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "BGP not initialized"); return -1; }
if (node_id == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: cm_open node_id=0 — rejected grp=%016llx type=%d ch=%s",
(unsigned long long)group->group_id, group->group_type,
group->channel_id[0] ? group->channel_id : "-");
if (cb) cb(NULL, node_id, 0, CONN_EVENT_TIMEOUT, cb_arg);
return -1;
}
uint64_t gid = group->group_id;
struct TOPO_GROUP_NODE* target = topo_node_find_by_id(group, node_id);
int loaded_from_db = 0;
if (!target && mgr->instance->topo_sqlite_db && mgr->instance->topo_groups) {
struct TOPO_NODE* ni = topo_node_sqlite_node_load(mgr->instance->topo_sqlite_db, mgr->instance->topo_groups, node_id);
if (ni) {
ni = topo_node_registry_store(mgr->instance->topo_groups, ni);
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_GROUP_NODE));
if (qe) {
target = (struct TOPO_GROUP_NODE*)qe;
memset((uint8_t*)target + sizeof(struct ll_entry), 0, sizeof(*target) - sizeof(struct ll_entry));
target->node_id = node_id; target->conn_mgr_type = CONN_TYPE_NONE;
queue_data_put_with_index(group->nodes, &target->ll); loaded_from_db = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: node 0x%016llx loaded from DB", (unsigned long long)node_id);
} else DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: failed to alloc TOPO_GROUP_NODE for DB node 0x%016llx", (unsigned long long)node_id);
}
}
if (!target) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node 0x%016llx skipped — not in routing, no active connections", (unsigned long long)node_id); return -1; }
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id);
if (entry && entry->state == CONN_MGR_STATE_CONNECTED) {
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) return -1; if (out_handle) *out_handle = h;
h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
}
if (entry && entry->state == CONN_MGR_STATE_CONNECTING) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node 0x%016llx already connecting, adding handle", (unsigned long long)node_id);
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) return -1; if (out_handle) *out_handle = h; return 0;
}
struct ETCP_CONN* existing = topo_group_find_conn_for_node(group, node_id);
if (!existing) existing = instance_find_conn(mgr->instance, node_id);
if (existing && existing->peer_node_id == node_id && existing->links) {
struct ETCP_LINK* l = existing->links;
while (l) { if (l->link_state == 3 && l->initialized) break; l = l->next; }
if (l) {
entry = cm_ensure_entry(mgr, node_id);
if (!entry) return -1;
entry->state = CONN_MGR_STATE_CONNECTED; entry->conn_type = CONN_TYPE_DIRECT;
entry->idle_timeout_ms = idle_timeout_ms; entry->last_traffic_tb = get_time_tb();
{ struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, node_id);
struct TOPO_NODE ni_e; memset(&ni_e, 0, sizeof(ni_e)); ni_e.node_id = node_id;
if (ni) memcpy(ni_e.public_key, ni->public_key, SC_PUBKEY_SIZE);
node_conn_direct_open_node(mgr->instance, node_id, cm_ncd_callback, entry, &entry->ncd_handle, &ni_e, NULL); }
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) { cm_entry_cleanup(entry); return -1; } if (out_handle) *out_handle = h;
cm_update_nodeinfo(entry); h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
}
}
entry = cm_ensure_entry(mgr, node_id);
if (!entry) return -1;
entry->state = CONN_MGR_STATE_CONNECTING;
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) { cm_entry_cleanup(entry); return -1; } if (out_handle) *out_handle = h;
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->db_loaded = (uint8_t)loaded_from_db;
if (cm_has_local_addr(mgr, group->local_node) && cm_has_local_addr(mgr, target))
{ 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 0;
}
void conn_mgr_close(struct CONN_MGR_HANDLE* h) {
if (!h) return;
struct CONN_MGR_ENTRY* entry = h->entry;
if (!entry) { u_free(h); return; }
/* отменяем отложенный UP-коллбэк (cm_deliver_up_cb), иначе он дёрнется
* с освобождённым handle (use-after-free при быстром teardown соединения) */
if (h->deliver_up_token && entry->mgr && entry->mgr->instance) {
uasync_call_soon_cancel(entry->mgr->instance->ua, h->deliver_up_token);
h->deliver_up_token = NULL;
}
{ struct CONN_MGR_HANDLE** pp = &entry->handles;
while (*pp) { if (*pp == h) { *pp = h->next; break; } pp = &(*pp)->next; }
}
if (entry->handles) { u_free(h); return; }
if (entry->state == CONN_MGR_STATE_CONNECTED && (entry->conn_type == CONN_TYPE_DIRECT || entry->conn_type == CONN_TYPE_REVERSE)) {
/* DISCONNECT через etcp_router — conn_mgr уровень */
struct CM_DISCONNECT pkt; memset(&pkt, 0, sizeof(pkt));
pkt.cmd = ETCP_RT_ID_CONN_MGR; pkt.subcmd = CM_SUBCMD_DISCONNECT; pkt.node_id = entry->node_id;
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = u_malloc(sizeof(pkt)); memcpy(qe->dgram, &pkt, sizeof(pkt)); qe->len = sizeof(pkt);
etcp_route_send(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, qe, 1); }
/* NCD CLOSE — транспортный уровень (отправит CLOSE/KEEP_ALIVE через node_conn_direct) */
}
cm_entry_cleanup(entry); u_free(h);
}
struct ETCP_CONN* conn_mgr_get_conn(struct CONN_MGR_HANDLE* h) {
if (!h || !h->entry || !h->entry->ncd_handle) return NULL;
return node_conn_direct_get_conn(h->entry->ncd_handle);
}
int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e) {
if (!h || !e) return -1;
struct CONN_MGR_ENTRY* entry = h->entry;
if (!entry || entry->state != CONN_MGR_STATE_CONNECTED) { queue_entry_free(e); return -1; }
entry->last_traffic_tb = get_time_tb();
struct ETCP_ROUTER_CONN* r = etcp_router_conn_get(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, ETCP_RT_ID_CONN_MGR);
if (!r) { queue_entry_free(e); return -1; }
int ret = etcp_router_conn_send(r, e->dgram, e->len); queue_entry_free(e); return ret;
}
/* ═══════ локальное сканирование ═══════ */
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_ping_cb_impl(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len);
struct cm_ping_ctx { struct CONN_MGR_ENTRY* entry; struct sockaddr_storage addr; struct ETCP_SOCKET* sock; uint8_t phase, attempt; };
/* LAN broadcast ping для обнаружения узла в локальной сети. Если ответит —
* соединение готово быстрее чем через интернет, без NAT. Запускается
* параллельно с DIRECT фазой. */
void cm_start_local_scan(struct CONN_MGR_ENTRY* entry) {
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* target = topo_node_find_by_id(group, entry->node_id);
struct TOPO_NODE* ni = target ? topo_node_registry_find(entry->mgr->instance->topo_groups, target->node_id) : NULL;
if (!target || !ni) { entry->local_scan_state = CM_TRY_FAILED; return; }
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = ni->v4_sock_meta; m; m = m->next) {
if (m->id != a->socket_id) continue;
if (m->config_type != CFG_SERVER_TYPE_PRIVATE && m->config_type != CFG_SERVER_TYPE_LOCAL) continue;
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) {
if ((s->type == CFG_SERVER_TYPE_PRIVATE || s->type == CFG_SERVER_TYPE_LOCAL) && s->local_addr.ss_family == AF_INET) {
struct cm_ping_ctx* ctx = u_calloc(1, sizeof(struct cm_ping_ctx));
if (!ctx) continue;
ctx->entry = entry; ctx->sock = s; ctx->phase = 0; ctx->attempt = 0;
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);
memcpy(&ctx->addr, &sin, sizeof(sin));
etcp_send_ping_to_socket(entry->mgr->instance, s, ni->public_key, &ctx->addr, CONN_MGR_LOCAL_SCAN_TIMEOUT_MS, cm_ping_cb, ctx, NULL, 0, 0);
return;
}
s = s->next;
}
}
}
entry->local_scan_state = CM_TRY_FAILED;
}
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)
{ cm_ping_cb_impl(success, rtt, arg, nonce, resp_data, resp_data_len); }
static void cm_ping_cb_impl(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_event(entry, CONN_EVENT_UP); u_free(ctx); return;
}
if (++ctx->attempt < CONN_MGR_LOCAL_SCAN_ATTEMPTS && ctx->phase == 0) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(entry->mgr->group, entry->node_id);
struct TOPO_NODE* pub_ni = nq ? topo_node_registry_find(entry->mgr->instance->topo_groups, nq->node_id) : NULL;
etcp_send_ping_to_socket(entry->mgr->instance, ctx->sock, pub_ni ? pub_ni->public_key : NULL,
&ctx->addr, CONN_MGR_LOCAL_SCAN_TIMEOUT_MS, cm_ping_cb, ctx, NULL, 0, 0);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: local scan failed for 0x%016llx after %d attempts",
(unsigned long long)entry->node_id, ctx->attempt);
entry->local_scan_state = CM_TRY_FAILED; u_free(ctx);
}
/* ═══════ DIRECT фаза (NCD + ручная NAT-фильтрация линков) ═══════ */
/* Добавляет IPv4/IPv6 линки к conn. Единый цикл для адресов из BGP и из базы:
* протокол из a->protocol (TCP — без NAT-фильтра, UDP — NAT-фильтрация),
* приоритет NAT > INTERFACE. Если ни одного линка не создалось — закрывает NCD
* и переходит к REVERSE/INDIRECT (для db_loaded — cleanup + TIMEOUT). */
static void cm_direct_add_links(struct ETCP_CONN* conn, struct CONN_MGR_ENTRY* entry,
struct TOPO_NODE* ni, int db_loaded) {
int any = 0, v4_cnt = 0, v6_cnt = 0, v4_tcp = 0, v6_tcp = 0;
/* v4 — единый цикл для адресов из BGP и из базы: протокол из a->protocol,
* приоритет NAT > INTERFACE, NAT-фильтрация только для UDP */
uint8_t prio[] = {TOPO_ADDR_NAT, TOPO_ADDR_INTERFACE};
for (int pr = 0; pr < 2; pr++) {
for (const struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) {
if (a->type != prio[pr]) 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 (a->protocol & TOPO_PROTO_TCP) {
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (s->is_tcp && s->local_addr.ss_family == AF_INET && s->type != CFG_SERVER_TYPE_PRIVATE) {
struct ETCP_LINK *tlink = etcp_link_new(conn, s, &sa, 0);
if (tlink) { tlink->is_tcp = 1; etcp_tcp_link_start_connect(tlink, &sa, a->port); any = 1; v4_tcp++; v4_cnt++; }
} s = s->next; }
}
if (a->protocol & TOPO_PROTO_UDP) {
for (const struct TOPO_SOCKMETA4* m = ni->v4_sock_meta; m; m = m->next) {
if (m->id != a->socket_id) continue;
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (cm_nat_compatible(s, m->config_type, m->nat_type) && s->local_addr.ss_family == AF_INET) {
if (etcp_link_new(conn, s, &sa, 0)) { any = 1; v4_cnt++; } } s = s->next; }
break;
}
}
}
}
/* v6 */
for (const struct TOPO_ADDR6* a6=ni->v6_addrs;a6;a6=a6->next) {
int is_ll = (a6->addr[0] == 0xfe && (a6->addr[1] & 0xc0) == 0x80);
if (a6->protocol & TOPO_PROTO_TCP) {
struct sockaddr_in6 sin6; memset(&sin6,0,sizeof(sin6)); sin6.sin6_family=AF_INET6;
memcpy(&sin6.sin6_addr,a6->addr,16); sin6.sin6_port=htons(a6->port);
if (is_ll) { struct ETCP_SOCKET* sv=entry->mgr->instance->etcp_sockets;
while(sv){if(sv->local_addr.ss_family==AF_INET6&&sv->netif_index)
{sin6.sin6_scope_id=sv->netif_index;break;} sv=sv->next;}
}
struct sockaddr_storage sa; memcpy(&sa,&sin6,sizeof(sin6));
{ struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (s->is_tcp && s->local_addr.ss_family == AF_INET6 && s->type != CFG_SERVER_TYPE_PRIVATE) {
struct ETCP_LINK *tlink = etcp_link_new(conn, s, &sa, 0);
if (tlink) { tlink->is_tcp = 1; etcp_tcp_link_start_connect(tlink, &sa, a6->port); any = 1; v6_tcp++; v6_cnt++; }
} s = s->next; }
}
}
if (a6->protocol & TOPO_PROTO_UDP) {
uint8_t tc=cm_classify_v6_addr(a6->addr); if(tc==CM_V6_OTH) continue;
struct ETCP_SOCKET* s=entry->mgr->instance->etcp_sockets;
while(s){uint8_t sc=cm_sock_v6_classify(s);if(sc==CM_V6_ANY||sc==tc){cm_add_v6_link(conn,a6->addr,a6->port,s);any=1; v6_cnt++;} s=s->next;}
}
}
if (any) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: direct_add_links node=0x%016llx db=%d v4=%d/tcp=%d v6=%d/tcp=%d conn=[%s]",
(unsigned long long)entry->node_id, db_loaded, v4_cnt, v4_tcp, v6_cnt, v6_tcp, conn->log_name);
return;
}
node_conn_direct_close(entry->ncd_handle); entry->ncd_handle=NULL;
if (db_loaded) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: db_node 0x%016llx direct failed (no compatible addr/socket)", (unsigned long long)entry->node_id);
cm_cleanup_db_node(entry); cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
} else {
struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; }
cm_start_phase_indirect(entry);
}
}
/* DIRECT фаза: открывает NCD-соединение с пустым ni (0 линков из NCD), затем
* вручную добавляет NAT-фильтрованные линки через cm_direct_add_links.
* NCD управляет таймаутом (2с). При успехе — cm_ncd_callback доставит UP,
* при таймауте — переключится на REVERSE или INDIRECT. */
void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
if (entry->main_connect_state == CM_TRY_OK || entry->local_scan_state == CM_TRY_OK) return;
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* target = topo_node_find_by_id(group, entry->node_id);
struct TOPO_NODE* ni = target ? topo_node_registry_find(entry->mgr->instance->topo_groups, target->node_id) : NULL;
if (!target || !ni) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
struct TOPO_NODE ni_e; memset(&ni_e, 0, sizeof(ni_e)); ni_e.node_id = target->node_id;
memcpy(ni_e.public_key, ni->public_key, SC_PUBKEY_SIZE);
int r = node_conn_direct_open_node(entry->mgr->instance, entry->node_id, cm_ncd_callback, entry,
&entry->ncd_handle, &ni_e, NULL);
if (r == NCD_ERR) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
cm_direct_add_links(node_conn_direct_get_conn(entry->ncd_handle), entry, ni, entry->db_loaded);
}
/* ═══════ REVERSE фаза ═══════ */
void cm_reverse_init_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)event;
struct cm_reverse_pending* rp = (struct cm_reverse_pending*)arg;
if (!rp) return; if (!conn) { u_free(rp); return; }
struct CONN_MGR_ENTRY* entry = rp->entry;
if (!entry) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse ready for 0x%016llx (no entry, freeing)", (unsigned long long)conn->peer_node_id); u_free(rp); return; }
if (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_event(entry, CONN_EVENT_UP); u_free(rp);
}
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);
}
/* REVERSE фаза: у нас прямой IP, у цели нет. Отправляем DIRECT_REQ с НАШИМИ
* адресами цели через BGP-маршрут ("подключись ко мне"). Ждём обратного INIT.
* Таймаут 15с — при провале переходим к INDIRECT. */
void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* local = group->local_node;
struct TOPO_NODE* lni = local ? topo_node_registry_find(entry->mgr->instance->topo_groups, local->node_id) : NULL;
if (!local || !lni) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
uint32_t req_id = ++entry->mgr->next_request_id; entry->main.request_id = req_id; entry->main.phase = 2;
uint8_t dc = 0; struct { uint8_t t; uint8_t ip[4]; uint16_t port; uint8_t sid; } addrs[8];
for (const struct TOPO_ADDR4* a = lni->v4_addrs; a && dc < 8; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = lni->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))
{ addrs[dc]=(typeof(addrs[0])){a->type}; memcpy(addrs[dc].ip,a->addr,4); addrs[dc].port=htons(a->port); addrs[dc].sid=a->socket_id; dc++; break; }
}
}
if (!dc) { cm_start_phase_indirect(entry); return; }
size_t sz = CM_DIRECT_REQ_H_SIZE + (size_t)dc * 8; uint8_t* pkt = u_malloc(sz);
if (!pkt) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
struct CM_DIRECT_REQ* req = (struct CM_DIRECT_REQ*)pkt; memset(req,0,sizeof(*req));
req->cmd=ETCP_RT_ID_CONN_MGR; req->subcmd=CM_SUBCMD_DIRECT_REQ; req->request_id=req_id; req->addr_count=dc;
uint8_t* p = pkt + CM_DIRECT_REQ_H_SIZE;
for (uint8_t i=0;i<dc;i++) { *p++=addrs[i].t; memcpy(p,addrs[i].ip,4);p+=4; memcpy(p,&addrs[i].port,2);p+=2; *p++=addrs[i].sid; }
struct ll_entry* qe = queue_entry_new(0);
if (!qe) { u_free(pkt); cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
qe->dgram=pkt; qe->len=(uint16_t)sz; etcp_route_send(entry->mgr->instance, group->group_id, entry->node_id, qe, 1);
struct cm_reverse_pending* rp = u_calloc(1, sizeof(*rp));
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, dc);
}
/* ═══════ обработка DIRECT_REQ (REVERSE входящие) ═══════ */
/* Принимающая сторона REVERSE: получили DIRECT_REQ от инициатора — открываем
* NCD-соединение к нему, добавляем линки по адресам из запроса. При INIT
* вызывается cm_reverse_init_cb. */
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 CM_DIRECT_REQ* req = (struct CM_DIRECT_REQ*)data;
if (len < CM_DIRECT_REQ_H_SIZE + (size_t)req->addr_count * 8) { 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);
uint64_t src = conn->peer_node_id;
struct TOPO_NODE* src_ni = topo_node_registry_find(mgr->instance->topo_groups, src); if (!src_ni) return;
struct TOPO_NODE nbuf; memset(&nbuf,0,sizeof(nbuf)); nbuf.node_id=src; memcpy(nbuf.public_key,src_ni->public_key,SC_PUBKEY_SIZE);
struct NODE_CONN_DIRECT* ncd_h = NULL;
if (node_conn_direct_open_node(mgr->instance, src, NULL, NULL, &ncd_h, &nbuf, NULL) == NCD_ERR) return;
struct ETCP_CONN* newc = node_conn_direct_get_conn(ncd_h); if (!newc) { node_conn_direct_close(ncd_h); return; }
struct cm_reverse_pending* rp = u_calloc(1,sizeof(*rp));
if (!rp) { node_conn_direct_close(ncd_h); return; }
rp->request_id=req->request_id; rp->next=mgr->reverse_pending; mgr->reverse_pending=rp;
uint8_t* p = (uint8_t*)data + CM_DIRECT_REQ_H_SIZE;
for (uint8_t i=0;i<req->addr_count;i++) {
uint8_t type=*p++; (void)type; uint8_t ip[4]; memcpy(ip,p,4);p+=4; uint16_t port; memcpy(&port,p,2);p+=2; uint8_t sid=*p++; (void)sid;
struct ETCP_SOCKET* s = conn->instance->etcp_sockets;
while (s) { if ((s->type==CFG_SERVER_TYPE_PUBLIC||s->type==CFG_SERVER_TYPE_UNKNOWN)&&s->local_addr.ss_family==AF_INET)
{ cm_add_v4_link(newc, ip, port, s); break; } s=s->next; }
}
etcp_conn_add_cbk(newc, cm_reverse_init_cb, rp, ETCP_CBK_EVENT_INIT);
}
/* ═══════ disconnect / recv ═══════ */
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_cleanup(entry);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: remote disconnect from 0x%016llx", (unsigned long long)node_id);
}
void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || !entry->dgram || entry->len < 3) return;
uint8_t* d = entry->dgram + 1; size_t len = entry->len - 1; uint8_t sub = d[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 (sub) {
case CM_SUBCMD_DIRECT_REQ: cm_handle_direct_req(conn, d, len); break;
case CM_SUBCMD_DIRECT_RESP: break;
case CM_SUBCMD_INTERM_EXCHANGE_REQ: if (len >= CM_EXCHANGE_REQ_SIZE) cm_handle_interm_exchange_req(conn, (struct CM_EXCHANGE_REQ*)d); break;
case CM_SUBCMD_INTERM_EXCHANGE_RESP: if (mgr) cm_handle_interm_exchange_resp(mgr, d, len); break;
case CM_SUBCMD_INTERM_SELECTED: if (mgr) cm_handle_interm_selected(mgr, d, len); break;
case CM_SUBCMD_DISCONNECT: if (len >= CM_DISCONNECT_SIZE && mgr) cm_handle_disconnect(mgr, ((struct CM_DISCONNECT*)d)->node_id); break;
}
}
int conn_mgr_open(struct UTUN_INSTANCE* inst,
uint64_t group_id,
uint64_t node_id,
conn_mgr_cb_t cb, void* cb_arg,
struct CONN_MGR_HANDLE** out_handle) {
if (out_handle) *out_handle = NULL;
if (!inst || !inst->topo_groups || !cb) return -1;
struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, group_id);
if (!group) {
char ch[32]; snprintf(ch, sizeof(ch), "%llu", (unsigned long long)group_id);
group = topo_groups_create_group(inst->topo_groups, group_id, TOPO_GROUP_TYPE_CHAT, ch);
if (!group) return -1;
}
struct CONN_MGR* mgr = group->conn_mgr;
if (!mgr) return -1;
return cm_open(mgr, node_id, 0, cb, cb_arg, out_handle);
}