Browse Source

conn_mgr: унифицировать адреса из BGP и базы, TCP sock_meta, протокол в DIRECT

topo_upd
evgeny 2 months ago
parent
commit
f550c80592
  1. 53
      src/routing_layer/conn_mgr_core.c
  2. 1
      src/routing_layer/conn_mgr_priv.h
  3. 2
      src/routing_layer/nat_detection.c
  4. 24
      src/routing_layer/topo_node.c
  5. 2
      src/transport_layer/etcp_keepalive.c
  6. 6
      src/transport_layer/stcp_link.c
  7. 2
      tests/test_nat_detection.c

53
src/routing_layer/conn_mgr_core.c

@ -55,7 +55,8 @@ static struct CONN_MGR_HANDLE* cm_invite_handle_new(conn_mgr_cb_t cb, void* cb_a
static void cm_deliver_up_cb(void* arg) {
struct CONN_MGR_HANDLE* h = (struct CONN_MGR_HANDLE*)arg;
if (h->cb) h->cb(h, h->node_id, h->entry->mgr->group->group_id, CONN_EVENT_UP, h->cb_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. Одно соединение — много слушателей:
@ -88,7 +89,7 @@ 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 (ot == CFG_SERVER_TYPE_NAT && our->nat_type == NAT_VERIFIED_EIM && tt == CFG_SERVER_TYPE_PUBLIC) 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;
@ -387,7 +388,7 @@ int conn_mgr_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_
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;
uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
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);
@ -411,7 +412,7 @@ int conn_mgr_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_
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); uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
cm_update_nodeinfo(entry); h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
}
}
@ -433,6 +434,12 @@ 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; }
}
@ -530,15 +537,19 @@ static void cm_ping_cb_impl(int success, uint16_t rtt, void* arg, uint64_t nonce
/* ═══════ DIRECT фаза (NCD + ручная NAT-фильтрация линков) ═══════ */
/* Добавляет IPv4/IPv6 линки к conn. db_loaded: все сокеты без фильтра.
* Обычный путь: NAT-фильтрация (cm_nat_compatible) + приоритет NAT > INTERFACE.
* Если ни одного линка не создалось — закрывает NCD и переходит к REVERSE/INDIRECT. */
/* Добавляет 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 */
if (db_loaded) {
/* 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));
@ -550,23 +561,13 @@ static void cm_direct_add_links(struct ETCP_CONN* conn, struct CONN_MGR_ENTRY* e
} s = s->next; }
}
if (a->protocol & TOPO_PROTO_UDP) {
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (s->local_addr.ss_family == AF_INET) {
if (etcp_link_new(conn,s,&sa,0)) any=1; } s=s->next; }
}
}
} else {
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;
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){
struct sockaddr_in sin; memset(&sin,0,sizeof(sin)); sin.sin_family=AF_INET;
memcpy(&sin.sin_addr.s_addr,a->addr,4); sin.sin_port=htons(a->port);
struct sockaddr_storage sa; memcpy(&sa,&sin,sizeof(sin));
if(etcp_link_new(conn,s,&sa,0)) any=1;} s=s->next;}
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;
}
}
}
}

1
src/routing_layer/conn_mgr_priv.h

@ -146,6 +146,7 @@ struct cm_invite_pending {
struct CONN_MGR_HANDLE {
uint64_t node_id; uint64_t group_id; conn_mgr_cb_t cb; void* cb_arg;
struct CONN_MGR_ENTRY* entry; struct CONN_MGR_HANDLE* next;
void* deliver_up_token; /* токен отложенного cm_deliver_up_cb (uasync_call_soon), NULL если не постился/уже сработал */
};
struct CONN_MGR_ENTRY {

2
src/routing_layer/nat_detection.c

@ -180,7 +180,7 @@ static void nat_detection_handle_nat_info(struct NAT_DETECTION* nd, struct TOPO_
if (a->type == TOPO_ADDR_INTERFACE && a->socket_id == socket_id) {
uint32_t a_ip; memcpy(&a_ip, a->addr, 4);
if (a_ip == nat_ip && a->port == nat_port) {
a->type = TOPO_ADDR_NAT; a->socket_id = socket_id | 1;
a->type = TOPO_ADDR_NAT; a->socket_id = socket_id;
dirty = 1;
ni->ver = (ni->ver % 255) + 1;
group->local_node->last_ver = ni->ver;

24
src/routing_layer/topo_node.c

@ -705,6 +705,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
int tcp4_count = 0, tcp6_count = 0;
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock) {
if (e_sock->is_tcp) { e_sock = e_sock->next; continue; }
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; }
if (e_sock->local_addr.ss_family == AF_INET) { sock_count++; addr_count++; }
else if (e_sock->local_addr.ss_family == AF_INET6) { sock6_count++; addr6_count++; }
@ -713,8 +714,8 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
{ struct ETCP_SOCKET* ts = instance->etcp_sockets;
while (ts) { if (!ts->is_tcp) { ts = ts->next; continue; }
if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; }
if (ts->interface_addr.ss_family == AF_INET || ts->local_addr.ss_family == AF_INET) { tcp4_count++; addr_count++; }
else if (ts->interface_addr.ss_family == AF_INET6 || ts->local_addr.ss_family == AF_INET6) { tcp6_count++; addr6_count++; }
if (ts->interface_addr.ss_family == AF_INET || ts->local_addr.ss_family == AF_INET) { tcp4_count++; sock_count++; addr_count++; }
else if (ts->interface_addr.ss_family == AF_INET6 || ts->local_addr.ss_family == AF_INET6) { tcp6_count++; sock6_count++; addr6_count++; }
ts = ts->next; }
}
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "counts: udp_v4s=%d udp_v4a=%d udp_v6s=%d udp_v6a=%d tcp4=%d tcp6=%d v4a_total=%d v6a_total=%d",
@ -794,6 +795,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: etcp_iter[%d] e_sock=%p type=%d fam=%d next=%p",
etcp_iter_cnt, (void*)e_sock, e_sock->type, e_sock->local_addr.ss_family, (void*)e_sock->next);
etcp_iter_cnt++;
if (e_sock->is_tcp) { e_sock = e_sock->next; continue; }
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; }
if (e_sock->local_addr.ss_family == AF_INET) {
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool);
@ -807,7 +809,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
{ struct TOPO_ADDR4* a = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool);
if (nat_verified) {
memcpy(a->addr, &nat_sin->sin_addr.s_addr, 4); a->port = ntohs(nat_sin->sin_port);
a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id | 1;
a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id;
} else {
int use_local = (local_sin->sin_addr.s_addr != 0);
memcpy(a->addr, use_local ? &local_sin->sin_addr.s_addr : &if_sin->sin_addr.s_addr, 4);
@ -833,16 +835,23 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
while (ts) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: tcp_iter[%d] ts=%p type=%d fam=%d next=%p",
tcp_iter_cnt, (void*)ts, ts->type, ts->local_addr.ss_family, (void*)ts->next);
tcp_iter_cnt++; if (!ts->is_tcp) { ts = ts->next; continue; }
if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; }
struct sockaddr_storage* addr = ts->interface_addr.ss_family ? &ts->interface_addr : &ts->local_addr;
if (!addr || !addr->ss_family) { ts = ts->next; continue; }
if (addr->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)addr;
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool);
sm->id = ts->sock_id; sm->config_type = ts->type; sm->nat_type = ts->nat_type;
sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
{ struct TOPO_ADDR4* a = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool);
memcpy(a->addr, &sin->sin_addr.s_addr, 4); a->port = ntohs(sin->sin_port);
a->type = TOPO_ADDR_INTERFACE; a->socket_id = ts->sock_id; a->protocol = TOPO_PROTO_TCP;
a->next = ni->v4_addrs; ni->v4_addrs = a; }
} else if (addr->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)addr;
{ struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(group->instance->topo_groups->v6_sock_meta_pool);
sm6->id = ts->sock_id; sm6->config_type = ts->type; sm6->nat_type = ts->nat_type;
sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; }
{ struct TOPO_ADDR6* a6 = memory_pool_alloc(group->instance->topo_groups->v6_addr_pool);
memcpy(a6->addr, &sin6->sin6_addr, 16); a6->port = ntohs(sin6->sin6_port);
a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = ts->sock_id; a6->protocol = TOPO_PROTO_TCP;
@ -902,6 +911,7 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock) {
if (e_sock->is_tcp) { e_sock = e_sock->next; continue; }
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; }
if (e_sock->local_addr.ss_family == AF_INET) {
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(instance->topo_groups->v4_sock_meta_pool);
@ -914,7 +924,7 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
{ struct TOPO_ADDR4* a = memory_pool_alloc(instance->topo_groups->v4_addr_pool);
if (nat_verified) {
memcpy(a->addr, &nat_sin->sin_addr.s_addr, 4); a->port = ntohs(nat_sin->sin_port);
a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id | 1;
a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id;
} else {
int use_local = (local_sin->sin_addr.s_addr != 0);
memcpy(a->addr, use_local ? &local_sin->sin_addr.s_addr : &if_sin->sin_addr.s_addr, 4);
@ -942,12 +952,18 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (!addr || !addr->ss_family) { ts = ts->next; continue; }
if (addr->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)addr;
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(instance->topo_groups->v4_sock_meta_pool);
sm->id = ts->sock_id; sm->config_type = ts->type; sm->nat_type = ts->nat_type;
sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
{ struct TOPO_ADDR4* a = memory_pool_alloc(instance->topo_groups->v4_addr_pool);
memcpy(a->addr, &sin->sin_addr.s_addr, 4); a->port = ntohs(sin->sin_port);
a->type = TOPO_ADDR_INTERFACE; a->socket_id = ts->sock_id; a->protocol = TOPO_PROTO_TCP;
a->next = ni->v4_addrs; ni->v4_addrs = a; }
} else if (addr->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)addr;
{ struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(instance->topo_groups->v6_sock_meta_pool);
sm6->id = ts->sock_id; sm6->config_type = ts->type; sm6->nat_type = ts->nat_type;
sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; }
{ struct TOPO_ADDR6* a6 = memory_pool_alloc(instance->topo_groups->v6_addr_pool);
memcpy(a6->addr, &sin6->sin6_addr, 16); a6->port = ntohs(sin6->sin6_port);
a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = ts->sock_id; a6->protocol = TOPO_PROTO_TCP;

2
src/transport_layer/etcp_keepalive.c

@ -115,8 +115,6 @@ static void keepalive_timer_cb(void* arg) {
uint64_t now = get_time_tb();
uint64_t timeout_units = (uint64_t)link->keepalive_timeout * 10; // ms -> 0.1ms units
uint64_t elapsed = now - link->last_recv_local_time;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] ka_check: recv=%d elapsed=%llu timeout=%llu",
link->etcp->log_name, link->recv_keepalive, (unsigned long long)elapsed, (unsigned long long)timeout_units);
if (elapsed > timeout_units) {
if (link->recv_keepalive != 0) {

6
src/transport_layer/stcp_link.c

@ -60,12 +60,6 @@ static void link_rx_cb(struct ll_queue *q, void *arg) {
if (!e) { queue_resume_callback(q); return; }
if (link->etcp_conn && link->etcp_link && link->inst) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[%s] stcp_rx: len=%u hdr=%02x%02x%02x sock=%p",
link->etcp_conn->log_name, e->len,
e->dgram && e->len>=3 ? e->dgram[0] : 0,
e->dgram && e->len>=3 ? e->dgram[1] : 0,
e->dgram && e->len>=3 ? e->dgram[2] : 0,
(void*)link->etcp_link->conn);
if (e->len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "stcp_link rx oversized pkt: %u > %u, closing link=%p", e->len, (unsigned)PACKET_DATA_SIZE, (void*)link);
queue_dgram_free(e); queue_entry_free(e);

2
tests/test_nat_detection.c

@ -378,7 +378,7 @@ int main(void) {
// Verify NAT addr in flat address list
int nat_found = 0;
for (const struct TOPO_ADDR4* a = lni->v4_addrs; a; a = a->next) {
if (a->type == TOPO_ADDR_NAT && a->socket_id == (sock_c1->sock_id | 1)) {
if (a->type == TOPO_ADDR_NAT && a->socket_id == sock_c1->sock_id) {
uint32_t addr_ip; memcpy(&addr_ip, a->addr, 4);
if (addr_ip != link_sc1->nat_ip || a->port != link_sc1->nat_port) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "FAIL: NAT addr mismatch in local_node: addr=%08x nat_ip=%08x port=%u nat_port=%u",

Loading…
Cancel
Save