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.
 
 
 
 
 
 

402 lines
19 KiB

#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 "topo_node.h"
#include "topo_group.h"
#include "route_connectivity.h"
#define CONN_MAX_SOCKET_CANDIDATES 8
struct conn_probe_ctx {
struct conn_probe_ctx* next; /* linked list in nq->connectivity.probe_list */
struct UTUN_INSTANCE* instance;
struct TOPO_NODEQ* 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;
}
if (a->ss_family == AF_INET6) {
const struct sockaddr_in6* sa = (const struct sockaddr_in6*)a;
const struct sockaddr_in6* sb = (const struct sockaddr_in6*)b;
if (memcmp(&sa->sin6_addr, &sb->sin6_addr, 16) != 0) return 1;
if (sa->sin6_port != sb->sin6_port) return 1;
return 0;
}
return memcmp(a, b, sizeof(struct sockaddr_storage));
}
// проверяет что два адреса (NODEINFO_IPV4_ADDR) не дубликаты по IP+port
static int addr_eq(const struct TOPO_ADDR4* 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;
int target_family = target_addr->ss_family;
int found = 0;
struct ETCP_SOCKET* e_sock;
if (target_family == AF_INET6) {
struct sockaddr_in6* target_sin6 = (struct sockaddr_in6*)target_addr;
// IPv6: prefer sockets in same /64 subnet
e_sock = instance->etcp_sockets;
while (e_sock && found < (int)max_count) {
if (e_sock->local_addr.ss_family == AF_INET6) {
struct sockaddr_in6* if_sin6 = (struct sockaddr_in6*)&e_sock->interface_addr;
if (memcmp(&if_sin6->sin6_addr, &target_sin6->sin6_addr, 8) == 0) {
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;
}
// Pass 2: all other IPv6 sockets
e_sock = instance->etcp_sockets;
while (e_sock && found < (int)max_count) {
if (e_sock->local_addr.ss_family != AF_INET6) { 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;
}
// IPv4 target — existing logic
uint32_t target_ip = ((const struct sockaddr_in*)target_addr)->sin_addr.s_addr;
e_sock = instance->etcp_sockets;
// Проход 1: точное совпадение подсети (для INTERFACE) или PUBLIC/NAT_VERIFIED (для NAT)
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 == TOPO_ADDR_INTERFACE) {
if ((e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL)
&& (if_ip & 0x00FFFFFF) == (target_ip & 0x00FFFFFF))
match = 1;
else if (e_sock->nat_type == NAT_VERIFIED_DIRECT || e_sock->nat_type == NAT_VERIFIED_EIM
|| e_sock->type == CFG_SERVER_TYPE_PUBLIC || e_sock->type == CFG_SERVER_TYPE_UNKNOWN)
match = 1;
} else {
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 || e_sock->type == CFG_SERVER_TYPE_UNKNOWN) 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 == TOPO_ADDR_INTERFACE) {
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL
|| e_sock->nat_type == NAT_VERIFIED_DIRECT || e_sock->nat_type == NAT_VERIFIED_EIM
|| e_sock->type == CFG_SERVER_TYPE_PUBLIC || e_sock->type == CFG_SERVER_TYPE_UNKNOWN)
out_sockets[found++] = e_sock;
} else {
if (e_sock->type != CFG_SERVER_TYPE_PRIVATE && e_sock->type != CFG_SERVER_TYPE_LOCAL)
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;
static const char* atype_names[] = {"INTERFACE", "NAT", "REAL"};
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "probe series start: socket=%s target=%s:%u type=%s",
sock->name,
ip_to_str(&((struct sockaddr_in*)&ctx->target_addr)->sin_addr, AF_INET).str,
(unsigned)ntohs(((struct sockaddr_in*)&ctx->target_addr)->sin_port),
atype_names[ctx->addr_type]);
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_DEBUG(DEBUG_CATEGORY_BGP, "probe ping: ok=%d rtt=%u sent=%d/%d socket=%s",
success, rtt, ctx->count_sent, ctx->count_total,
ctx->candidate_sockets[ctx->candidate_index]->name);
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) return;
if (ctx->nq) {
struct TOPO_CONNECTIVITY* c = &ctx->nq->connectivity;
struct conn_probe_ctx** p = (struct conn_probe_ctx**)&c->probe_list;
while (*p) { if (*p == ctx) { *p = ctx->next; break; } p = &(*p)->next; }
switch (ctx->addr_type) {
case TOPO_ADDR_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();
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;
case TOPO_ADDR_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;
}
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 TOPO_NODEQ* 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;
}
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;
int pend = 0;
// ---- IPv4 addresses ----
if (nq->node && nq->node->v4_addrs) {
uint32_t seen_ips[16]; uint16_t seen_ports[16]; int seen_count = 0;
for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) {
uint32_t ip; memcpy(&ip, a->addr, 4); uint16_t port = a->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*)&target; 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, a->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, a->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 = a->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;
ctx->next = (struct conn_probe_ctx*)nq->connectivity.probe_list; nq->connectivity.probe_list = ctx; pend++;
conn_probe_start_series(ctx);
}
}
// ---- IPv6 addresses ----
if (nq->node && nq->node->v6_addrs) {
uint8_t seen_ip6s[16][16]; uint16_t seen_ports6[16]; int seen6_count = 0;
for (const struct TOPO_ADDR6* a6 = nq->node->v6_addrs; a6; a6 = a6->next) {
if (a6->port == 0) continue;
int zero = 1; for (int z = 0; z < 16; z++) if (a6->addr[z] != 0) { zero = 0; break; } if (zero) continue;
int dup = 0; for (int j = 0; j < seen6_count; j++) { if (memcmp(seen_ip6s[j], a6->addr, 16) == 0 && seen_ports6[j] == a6->port) { dup = 1; break; } }
if (dup) continue;
if (seen6_count >= 16) break;
memcpy(seen_ip6s[seen6_count], a6->addr, 16); seen_ports6[seen6_count] = a6->port; seen6_count++;
struct sockaddr_storage target; memset(&target, 0, sizeof(target));
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&target; sin6->sin6_family = AF_INET6;
memcpy(&sin6->sin6_addr, a6->addr, 16); sin6->sin6_port = htons(a6->port);
struct ETCP_SOCKET* candidates[CONN_MAX_SOCKET_CANDIDATES];
int cand_count = conn_match_candidate_sockets(instance, a6->type, &target, candidates, CONN_MAX_SOCKET_CANDIDATES);
if (cand_count == 0) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe: no local sockets for IPv6 target type=%d", a6->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 = a6->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;
ctx->next = (struct conn_probe_ctx*)nq->connectivity.probe_list; nq->connectivity.probe_list = ctx; pend++;
conn_probe_start_series(ctx);
}
}
if (pend == 0) {
nq->connectivity.probe_status = PROBE_STATUS_DONE;
if (!nq->node->v4_addrs && !nq->node->v6_addrs) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "no addresses to probe for node %016llx", (unsigned long long)nq->node->node_id);
}
} else {
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 TOPO_NODEQ* nq) {
if (!instance || !nq) return;
nq->connectivity.probe_status = PROBE_STATUS_NONE;
nq->connectivity.pending_count = 0;
struct conn_probe_ctx* ctx = (struct conn_probe_ctx*)nq->connectivity.probe_list;
while (ctx) { ctx->nq = NULL; ctx = ctx->next; }
nq->connectivity.probe_list = NULL;
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) {
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
if (!instance || !group || !group->nodes) return;
int count = 0;
struct ll_entry* e = group->nodes->head;
while (e) {
struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)e;
nq->connectivity.probe_status = PROBE_STATUS_NONE;
nq->connectivity.pending_count = 0;
e = e->next;
count++;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe cancel ALL: %d nodes", count);
}