Browse Source
- Rename NAT types: OPEN->EIM, RESTRICTED->STRICT across all code
- Add flat addr list {type,ip,port,socket_id} to NODEINFO (NODEINFO_IPV4_ADDR)
- Add NODEINFO_IPV4_SOCKET_META for per-socket metadata
- Add NODE_CONNECTIVITY with per-type probe status and min_rtt
- New module: route_connectivity.c/h (probing engine with socket fallback)
- STRICT verified addresses excluded from NODEINFO broadcast
- Trigger probing on new/updated BGP node, cancel on remove/withdraw
- Update NAT type names in etcpmon GUI
congestion
14 changed files with 807 additions and 473 deletions
@ -0,0 +1,348 @@ |
|||||||
|
#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 "route_node.h" |
||||||
|
#include "route_bgp.h" |
||||||
|
#include "route_connectivity.h" |
||||||
|
|
||||||
|
#define CONN_MAX_SOCKET_CANDIDATES 8 |
||||||
|
|
||||||
|
struct conn_probe_ctx { |
||||||
|
struct UTUN_INSTANCE* instance; |
||||||
|
struct NODEINFO_Q* nq; |
||||||
|
uint8_t addr_type; // ADDR_TYPE_*
|
||||||
|
struct sockaddr_storage target_addr; |
||||||
|
uint8_t peer_pubkey[SC_PUBKEY_SIZE]; |
||||||
|
|
||||||
|
struct ETCP_SOCKET* candidate_sockets[CONN_MAX_SOCKET_CANDIDATES]; |
||||||
|
uint8_t candidate_count; |
||||||
|
uint8_t candidate_index; |
||||||
|
|
||||||
|
uint16_t best_across_sockets; // min RTT по всем сокетам
|
||||||
|
uint8_t count_total; // 3 на серию
|
||||||
|
uint8_t count_sent; |
||||||
|
uint8_t count_ok; |
||||||
|
uint16_t min_rtt; // min RTT в текущей серии
|
||||||
|
uint16_t timeout_ms; |
||||||
|
void* ping_timer; |
||||||
|
}; |
||||||
|
|
||||||
|
// ---- forward ----
|
||||||
|
static void conn_probe_single_cb(int success, uint16_t rtt, void* arg, |
||||||
|
uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len); |
||||||
|
static void conn_probe_finish(struct conn_probe_ctx* ctx, int ok); |
||||||
|
static void conn_probe_start_series(struct conn_probe_ctx* ctx); |
||||||
|
|
||||||
|
// ---- helpers ----
|
||||||
|
|
||||||
|
static int sock_addr_cmp(const struct sockaddr_storage* a, const struct sockaddr_storage* b) { |
||||||
|
if (a->ss_family != b->ss_family) return 1; |
||||||
|
if (a->ss_family == AF_INET) { |
||||||
|
const struct sockaddr_in* sa = (const struct sockaddr_in*)a; |
||||||
|
const struct sockaddr_in* sb = (const struct sockaddr_in*)b; |
||||||
|
if (sa->sin_addr.s_addr != sb->sin_addr.s_addr) return 1; |
||||||
|
if (sa->sin_port != sb->sin_port) return 1; |
||||||
|
return 0; |
||||||
|
} |
||||||
|
return memcmp(a, b, sizeof(struct sockaddr_storage)); |
||||||
|
} |
||||||
|
|
||||||
|
// проверяет что два адреса (NODEINFO_IPV4_ADDR) не дубликаты по IP+port
|
||||||
|
static int addr_eq(const struct NODEINFO_IPV4_ADDR* a, uint32_t ip, uint16_t port) { |
||||||
|
uint32_t a_ip; memcpy(&a_ip, a->addr, 4); |
||||||
|
return a_ip == ip && a->port == port; |
||||||
|
} |
||||||
|
|
||||||
|
// собирает список локальных сокетов-кандидатов для probing заданного адреса
|
||||||
|
// сортирует: лучшие (по совпадению подсети / типу NAT) первые
|
||||||
|
static int conn_match_candidate_sockets(struct UTUN_INSTANCE* instance, |
||||||
|
uint8_t addr_type, |
||||||
|
const struct sockaddr_storage* target_addr, |
||||||
|
struct ETCP_SOCKET** out_sockets, uint8_t max_count) { |
||||||
|
if (!instance || !target_addr || !out_sockets || max_count == 0) return 0; |
||||||
|
|
||||||
|
uint32_t target_ip = ((const struct sockaddr_in*)target_addr)->sin_addr.s_addr; |
||||||
|
int found = 0; |
||||||
|
struct ETCP_SOCKET* e_sock = instance->etcp_sockets; |
||||||
|
|
||||||
|
// Проход 1: точное совпадение подсети (для INTERFACE) или PUBLIC/NAT_VERIFIED (для NAT/REAL)
|
||||||
|
while (e_sock && found < (int)max_count) { |
||||||
|
if (e_sock->local_addr.ss_family != AF_INET) { e_sock = e_sock->next; continue; } |
||||||
|
int match = 0; |
||||||
|
struct sockaddr_in* if_sin = (struct sockaddr_in*)&e_sock->interface_addr; |
||||||
|
uint32_t if_ip = if_sin->sin_addr.s_addr; |
||||||
|
|
||||||
|
if (addr_type == ADDR_TYPE_INTERFACE) { |
||||||
|
// предпочитаем PRIVATE сокеты в той же /24 подсети
|
||||||
|
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE && (if_ip & 0x00FFFFFF) == (target_ip & 0x00FFFFFF)) |
||||||
|
match = 1; |
||||||
|
} else { |
||||||
|
// NAT или REAL: предпочитаем DIRECT > EIM > PUBLIC
|
||||||
|
if (e_sock->nat_type == NAT_VERIFIED_DIRECT) match = 1; |
||||||
|
else if (e_sock->nat_type == NAT_VERIFIED_EIM) match = 1; |
||||||
|
else if (e_sock->type == CFG_SERVER_TYPE_PUBLIC) match = 1; |
||||||
|
} |
||||||
|
|
||||||
|
if (match) { |
||||||
|
int dup = 0; |
||||||
|
for (int i = 0; i < found; i++) if (out_sockets[i] == e_sock) { dup = 1; break; } |
||||||
|
if (!dup) out_sockets[found++] = e_sock; |
||||||
|
} |
||||||
|
e_sock = e_sock->next; |
||||||
|
} |
||||||
|
|
||||||
|
// Проход 2: все остальные подходящие
|
||||||
|
e_sock = instance->etcp_sockets; |
||||||
|
while (e_sock && found < (int)max_count) { |
||||||
|
if (e_sock->local_addr.ss_family != AF_INET) { e_sock = e_sock->next; continue; } |
||||||
|
int dup = 0; |
||||||
|
for (int i = 0; i < found; i++) if (out_sockets[i] == e_sock) { dup = 1; break; } |
||||||
|
if (!dup) { |
||||||
|
if (addr_type == ADDR_TYPE_INTERFACE) { |
||||||
|
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE) out_sockets[found++] = e_sock; |
||||||
|
} else { |
||||||
|
// любой не-private
|
||||||
|
if (e_sock->type != CFG_SERVER_TYPE_PRIVATE) out_sockets[found++] = e_sock; |
||||||
|
} |
||||||
|
} |
||||||
|
e_sock = e_sock->next; |
||||||
|
} |
||||||
|
|
||||||
|
// Проход 3: остальные IPv4 (fallback)
|
||||||
|
e_sock = instance->etcp_sockets; |
||||||
|
while (e_sock && found < (int)max_count) { |
||||||
|
if (e_sock->local_addr.ss_family != AF_INET) { e_sock = e_sock->next; continue; } |
||||||
|
int dup = 0; |
||||||
|
for (int i = 0; i < found; i++) if (out_sockets[i] == e_sock) { dup = 1; break; } |
||||||
|
if (!dup) out_sockets[found++] = e_sock; |
||||||
|
e_sock = e_sock->next; |
||||||
|
} |
||||||
|
|
||||||
|
return found; |
||||||
|
} |
||||||
|
|
||||||
|
// ---- probe lifecycle ----
|
||||||
|
|
||||||
|
static void conn_probe_start_series(struct conn_probe_ctx* ctx) { |
||||||
|
if (ctx->candidate_index >= ctx->candidate_count) { |
||||||
|
// все сокеты перебраны
|
||||||
|
if (ctx->best_across_sockets != 65535) conn_probe_finish(ctx, 1); |
||||||
|
else conn_probe_finish(ctx, 0); |
||||||
|
return; |
||||||
|
} |
||||||
|
|
||||||
|
struct ETCP_SOCKET* sock = ctx->candidate_sockets[ctx->candidate_index]; |
||||||
|
ctx->count_sent = 0; |
||||||
|
ctx->count_ok = 0; |
||||||
|
ctx->min_rtt = 65535; |
||||||
|
|
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe series start: socket=%s (idx=%d/%d) addr_type=%d target=%s:%u", |
||||||
|
sock->name, ctx->candidate_index, ctx->candidate_count, ctx->addr_type, |
||||||
|
ip_to_str(&((struct sockaddr_in*)&ctx->target_addr)->sin_addr, AF_INET).str, |
||||||
|
(unsigned)ntohs(((struct sockaddr_in*)&ctx->target_addr)->sin_port)); |
||||||
|
|
||||||
|
int ret = etcp_send_ping_to_socket(ctx->instance, sock, ctx->peer_pubkey, |
||||||
|
&ctx->target_addr, ctx->timeout_ms, |
||||||
|
conn_probe_single_cb, ctx, NULL, 0); |
||||||
|
if (ret != 0) { |
||||||
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "probe: cannot start ping from socket %s", sock->name); |
||||||
|
ctx->count_sent = 3; // simulate full failure
|
||||||
|
ctx->candidate_index++; |
||||||
|
conn_probe_start_series(ctx); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void conn_probe_single_cb(int success, uint16_t rtt, void* arg, |
||||||
|
uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len) { |
||||||
|
(void)nonce; (void)resp_data; (void)resp_data_len; |
||||||
|
struct conn_probe_ctx* ctx = (struct conn_probe_ctx*)arg; |
||||||
|
if (!ctx) return; |
||||||
|
|
||||||
|
ctx->count_sent++; |
||||||
|
if (success) { |
||||||
|
ctx->count_ok++; |
||||||
|
if (rtt < ctx->min_rtt) ctx->min_rtt = rtt; |
||||||
|
} |
||||||
|
|
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe cb: success=%d rtt=%u sent=%d ok=%d min_rtt=%u", |
||||||
|
success, rtt, ctx->count_sent, ctx->count_ok, ctx->min_rtt); |
||||||
|
|
||||||
|
if (ctx->count_sent < ctx->count_total) { |
||||||
|
// продолжаем с тем же сокетом
|
||||||
|
struct ETCP_SOCKET* sock = ctx->candidate_sockets[ctx->candidate_index]; |
||||||
|
int ret = etcp_send_ping_to_socket(ctx->instance, sock, ctx->peer_pubkey, |
||||||
|
&ctx->target_addr, ctx->timeout_ms, |
||||||
|
conn_probe_single_cb, ctx, NULL, 0); |
||||||
|
if (ret != 0) { |
||||||
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "probe: cannot continue ping from socket %s", sock->name); |
||||||
|
ctx->count_sent = ctx->count_total; // force finish series
|
||||||
|
ctx->count_ok = 0; |
||||||
|
} else return; // следующий пинг отправлен, ждём callback
|
||||||
|
} |
||||||
|
|
||||||
|
// серия из 3 пингов завершена
|
||||||
|
if (ctx->count_ok > 0) { |
||||||
|
if (ctx->min_rtt < ctx->best_across_sockets) ctx->best_across_sockets = ctx->min_rtt; |
||||||
|
conn_probe_finish(ctx, 1); |
||||||
|
} else { |
||||||
|
// этот сокет не подошёл — пробуем следующий
|
||||||
|
ctx->candidate_index++; |
||||||
|
conn_probe_start_series(ctx); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void conn_probe_finish(struct conn_probe_ctx* ctx, int ok) { |
||||||
|
if (!ctx || !ctx->nq) return; |
||||||
|
struct NODE_CONNECTIVITY* c = &ctx->nq->connectivity; |
||||||
|
|
||||||
|
switch (ctx->addr_type) { |
||||||
|
case ADDR_TYPE_INTERFACE: |
||||||
|
c->interface_status = ok ? PROBE_RESULT_REACHABLE : (c->interface_status == PROBE_RESULT_UNKNOWN ? PROBE_RESULT_UNREACHABLE : c->interface_status); |
||||||
|
if (ok && ctx->best_across_sockets < c->interface_min_rtt) c->interface_min_rtt = ctx->best_across_sockets; |
||||||
|
c->interface_probe_time = get_time_tb(); |
||||||
|
break; |
||||||
|
case ADDR_TYPE_NAT: |
||||||
|
c->nat_status = ok ? PROBE_RESULT_REACHABLE : (c->nat_status == PROBE_RESULT_UNKNOWN ? PROBE_RESULT_UNREACHABLE : c->nat_status); |
||||||
|
if (ok && ctx->best_across_sockets < c->nat_min_rtt) c->nat_min_rtt = ctx->best_across_sockets; |
||||||
|
c->nat_probe_time = get_time_tb(); |
||||||
|
break; |
||||||
|
case ADDR_TYPE_REAL: |
||||||
|
c->real_status = ok ? PROBE_RESULT_REACHABLE : (c->real_status == PROBE_RESULT_UNKNOWN ? PROBE_RESULT_UNREACHABLE : c->real_status); |
||||||
|
if (ok && ctx->best_across_sockets < c->real_min_rtt) c->real_min_rtt = ctx->best_across_sockets; |
||||||
|
c->real_probe_time = get_time_tb(); |
||||||
|
break; |
||||||
|
} |
||||||
|
|
||||||
|
if (c->pending_count > 0) c->pending_count--; |
||||||
|
if (c->pending_count == 0) { |
||||||
|
c->probe_status = PROBE_STATUS_DONE; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "connectivity probe DONE for node %016llx: intf=%d nat=%d real=%d", |
||||||
|
(unsigned long long)ctx->nq->node.node_id, |
||||||
|
c->interface_status, c->nat_status, c->real_status); |
||||||
|
} |
||||||
|
|
||||||
|
u_free(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
// ---- public API ----
|
||||||
|
|
||||||
|
void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct NODEINFO_Q* nq) { |
||||||
|
if (!instance || !nq) return; |
||||||
|
if (nq->node.node_id == instance->node_id) return; // не пингуем себя
|
||||||
|
if (nq->connectivity.probe_status == PROBE_STATUS_IN_PROGRESS) { |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe already in progress for node %016llx", (unsigned long long)nq->node.node_id); |
||||||
|
return; |
||||||
|
} |
||||||
|
|
||||||
|
const struct NODEINFO_IPV4_ADDR* addrs; |
||||||
|
int addr_count = get_node_v4_addrs(nq, &addrs); |
||||||
|
if (addr_count <= 0) { |
||||||
|
nq->connectivity.probe_status = PROBE_STATUS_DONE; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "no addresses to probe for node %016llx", (unsigned long long)nq->node.node_id); |
||||||
|
return; |
||||||
|
} |
||||||
|
|
||||||
|
// дедупликация: массив уже проверенных (IP, port) пар
|
||||||
|
uint32_t seen_ips[16]; uint16_t seen_ports[16]; |
||||||
|
int seen_count = 0; |
||||||
|
int pend = 0; |
||||||
|
|
||||||
|
nq->connectivity.probe_status = PROBE_STATUS_IN_PROGRESS; |
||||||
|
nq->connectivity.probe_start_time = get_time_tb(); |
||||||
|
nq->connectivity.interface_status = PROBE_RESULT_UNKNOWN; |
||||||
|
nq->connectivity.nat_status = PROBE_RESULT_UNKNOWN; |
||||||
|
nq->connectivity.real_status = PROBE_RESULT_UNKNOWN; |
||||||
|
nq->connectivity.interface_min_rtt = 65535; |
||||||
|
nq->connectivity.nat_min_rtt = 65535; |
||||||
|
nq->connectivity.real_min_rtt = 65535; |
||||||
|
nq->connectivity.pending_count = 0; |
||||||
|
|
||||||
|
for (int i = 0; i < addr_count; i++) { |
||||||
|
uint32_t ip; memcpy(&ip, addrs[i].addr, 4); |
||||||
|
uint16_t port = addrs[i].port; |
||||||
|
if (ip == 0 || port == 0) continue; |
||||||
|
|
||||||
|
// дедупликация
|
||||||
|
int dup = 0; |
||||||
|
for (int j = 0; j < seen_count; j++) { |
||||||
|
if (seen_ips[j] == ip && seen_ports[j] == port) { dup = 1; break; } |
||||||
|
} |
||||||
|
if (dup) continue; |
||||||
|
if (seen_count >= 16) break; |
||||||
|
seen_ips[seen_count] = ip; seen_ports[seen_count] = port; seen_count++; |
||||||
|
|
||||||
|
// собрать целевой адрес
|
||||||
|
struct sockaddr_storage target; |
||||||
|
memset(&target, 0, sizeof(target)); |
||||||
|
struct sockaddr_in* sin = (struct sockaddr_in*)⌖ |
||||||
|
sin->sin_family = AF_INET; |
||||||
|
sin->sin_addr.s_addr = ip; |
||||||
|
sin->sin_port = htons(port); |
||||||
|
|
||||||
|
// собрать кандидатские локальные сокеты
|
||||||
|
struct ETCP_SOCKET* candidates[CONN_MAX_SOCKET_CANDIDATES]; |
||||||
|
int cand_count = conn_match_candidate_sockets(instance, addrs[i].type, &target, |
||||||
|
candidates, CONN_MAX_SOCKET_CANDIDATES); |
||||||
|
if (cand_count == 0) { |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe: no local sockets for target %s:%u type=%d", |
||||||
|
ip_to_str(&ip, AF_INET).str, port, addrs[i].type); |
||||||
|
continue; |
||||||
|
} |
||||||
|
|
||||||
|
struct conn_probe_ctx* ctx = u_calloc(1, sizeof(struct conn_probe_ctx)); |
||||||
|
if (!ctx) continue; |
||||||
|
ctx->instance = instance; |
||||||
|
ctx->nq = nq; |
||||||
|
ctx->addr_type = addrs[i].type; |
||||||
|
ctx->target_addr = target; |
||||||
|
memcpy(ctx->peer_pubkey, nq->node.public_key, SC_PUBKEY_SIZE); |
||||||
|
memcpy(ctx->candidate_sockets, candidates, cand_count * sizeof(struct ETCP_SOCKET*)); |
||||||
|
ctx->candidate_count = cand_count; |
||||||
|
ctx->candidate_index = 0; |
||||||
|
ctx->count_total = CONN_PROBE_COUNT; |
||||||
|
ctx->timeout_ms = CONN_PROBE_TIMEOUT_MS; |
||||||
|
ctx->best_across_sockets = 65535; |
||||||
|
|
||||||
|
pend++; |
||||||
|
conn_probe_start_series(ctx); |
||||||
|
} |
||||||
|
|
||||||
|
if (pend == 0) { |
||||||
|
nq->connectivity.probe_status = PROBE_STATUS_DONE; |
||||||
|
} else { |
||||||
|
if (nq->node.local_v4_addrs == 0) nq->connectivity.interface_status = PROBE_RESULT_UNREACHABLE; |
||||||
|
nq->connectivity.pending_count = pend; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "connectivity probe started for node %016llx: %d series pending", |
||||||
|
(unsigned long long)nq->node.node_id, pend); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
void route_connectivity_cancel_node(struct UTUN_INSTANCE* instance, struct NODEINFO_Q* nq) { |
||||||
|
if (!instance || !nq) return; |
||||||
|
nq->connectivity.probe_status = PROBE_STATUS_NONE; |
||||||
|
nq->connectivity.pending_count = 0; |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "connectivity probe cancelled for node %016llx", (unsigned long long)nq->node.node_id); |
||||||
|
} |
||||||
|
|
||||||
|
void route_connectivity_cancel_all(struct UTUN_INSTANCE* instance) { |
||||||
|
if (!instance || !instance->bgp || !instance->bgp->nodes) return; |
||||||
|
struct ll_entry* e = instance->bgp->nodes->head; |
||||||
|
while (e) { |
||||||
|
struct NODEINFO_Q* nq = (struct NODEINFO_Q*)e; |
||||||
|
nq->connectivity.probe_status = PROBE_STATUS_NONE; |
||||||
|
nq->connectivity.pending_count = 0; |
||||||
|
e = e->next; |
||||||
|
} |
||||||
|
} |
||||||
@ -0,0 +1,20 @@ |
|||||||
|
#ifndef ROUTE_CONNECTIVITY_H |
||||||
|
#define ROUTE_CONNECTIVITY_H |
||||||
|
|
||||||
|
#include <stdint.h> |
||||||
|
#include "route_node.h" |
||||||
|
|
||||||
|
struct UTUN_INSTANCE; |
||||||
|
|
||||||
|
// Запускает зондирование связности ко всем адресам удалённого узла
|
||||||
|
void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, |
||||||
|
struct NODEINFO_Q* nq); |
||||||
|
|
||||||
|
// Отменяет все pending пробы для узла (при удалении / withdraw)
|
||||||
|
void route_connectivity_cancel_node(struct UTUN_INSTANCE* instance, |
||||||
|
struct NODEINFO_Q* nq); |
||||||
|
|
||||||
|
// Отменяет все pending пробы для всех узлов (при destroy)
|
||||||
|
void route_connectivity_cancel_all(struct UTUN_INSTANCE* instance); |
||||||
|
|
||||||
|
#endif // ROUTE_CONNECTIVITY_H
|
||||||
Loading…
Reference in new issue