Browse Source

ping refactor

congestion
jeka 6 months ago
parent
commit
b2ed0317a6
  1. 15
      src/etcp_connections.c
  2. 2
      src/route_bgp.c
  3. 583
      src/route_ping.c
  4. 1
      src/route_ping.h
  5. 2
      tools/etcpmon/etcpmon_protocol.h

15
src/etcp_connections.c

@ -885,7 +885,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "sendto failed, sock_err=%d", socket_get_error()); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "sendto failed, sock_err=%d", socket_get_error());
dgram->link->send_errors++; errcode=4; goto es_err; dgram->link->send_errors++; errcode=4; goto es_err;
} else { } else {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "sendto succeeded, sent=%zd bytes to port %d", sent, ntohs(((struct sockaddr_in*)addr)->sin_port)); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "sendto succeeded, sent=%zd bytes to port %d fd=%d", sent, ntohs(((struct sockaddr_in*)addr)->sin_port), dgram->link->conn->fd);
} }
return (int)sent; return (int)sent;
es_err: es_err:
@ -918,18 +918,15 @@ static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, socket_t fd, sc_context_
} }
memcpy(enc_buf + enc_buf_len, dgram->data + len, dgram->noencrypt_len); memcpy(enc_buf + enc_buf_len, dgram->data + len, dgram->noencrypt_len);
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6); socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6);
int rnd = rand() % 100;
ssize_t sent = enc_buf_len + dgram->noencrypt_len; ssize_t sent = enc_buf_len + dgram->noencrypt_len;
if (loss_rate == 0 || rnd >= loss_rate) { // uint8_t* xaddr=&((struct sockaddr_in*)addr)->sin_addr;
sent = socket_sendto(fd, enc_buf, enc_buf_len + dgram->noencrypt_len, (struct sockaddr*)addr, addr_len); //xaddr[0]=192; xaddr[1]=168; xaddr[2]=10; xaddr[3]=1;
} else { sent = socket_sendto(fd, enc_buf, enc_buf_len + dgram->noencrypt_len, (struct sockaddr*)addr, addr_len);
DEBUG_WARN(DEBUG_CATEGORY_BGP, "ping packet dropped by loss_rate (rnd=%d, loss_rate=%d%%)", rnd, loss_rate);
}
if (sent < 0) { if (sent < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "sendto failed for ping, err=%d", socket_get_error()); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "sendto failed for ping, err=%d addr=%s fd=%d len=%d", socket_get_error(), sockaddr_storage_to_str(addr).str, fd, enc_buf_len + dgram->noencrypt_len);
return -1; return -1;
} }
DEBUG_TRACE(DEBUG_CATEGORY_BGP, "ping sendto succeeded to %s sent=%zd bytes", sockaddr_storage_to_str(addr).str, sent); DEBUG_TRACE(DEBUG_CATEGORY_BGP, "ping sendto succeeded to %s sent=%zd bytes fd=%d", sockaddr_storage_to_str(addr).str, sent, fd);
return (int)sent; return (int)sent;
} }

2
src/route_bgp.c

@ -285,7 +285,7 @@ void route_bgp_destroy(struct UTUN_INSTANCE* instance) {
etcp_unbind(instance, ETCP_ID_ROUTE_ENTRY); etcp_unbind(instance, ETCP_ID_ROUTE_ENTRY);
route_ping_destroy_pending(instance->bgp); // route_ping_destroy_pending(instance->bgp);
struct ll_entry* e; struct ll_entry* e;
while ((e = queue_data_get(instance->bgp->senders_list)) != NULL) { while ((e = queue_data_get(instance->bgp->senders_list)) != NULL) {

583
src/route_ping.c

@ -18,40 +18,21 @@
#include "route_bgp.h" #include "route_bgp.h"
#include "route_ping.h" #include "route_ping.h"
struct route_ping_sock_ctx { struct route_ping_series_ctx {
struct ETCP_SOCKET* local_sock;
uint16_t avg_rtt; // средний RTT в 0.1ms
uint8_t count_sent;
uint8_t count_ok;
};
struct route_ping_req {
struct ETCP_CONN* reply_conn; struct ETCP_CONN* reply_conn;
uint64_t request_id; uint64_t request_id;
struct NODEINFO_Q* target_node;
uint8_t socket_count;
uint8_t completed_count;
uint16_t timeout_ms;
uint8_t recv_ipv4[4]; // IP:port с которого получен PING_REQ (STUN-like)
uint16_t recv_port; // network byte order
struct route_ping_sock_ctx sockets[0];
};
struct route_ping_series {
struct route_ping_req* req;
struct route_ping_sock_ctx* sock_ctx;
struct ETCP_SOCKET* local_sock;
struct sockaddr_storage target_addr; struct sockaddr_storage target_addr;
uint8_t pubkey[SC_PUBKEY_SIZE]; uint8_t pubkey[SC_PUBKEY_SIZE];
uint8_t count_total; struct ETCP_SOCKET* local_sock;
uint8_t count_sent; uint8_t count_total;
uint8_t count_ok; uint8_t count_sent;
uint32_t rtt_sum; uint8_t count_ok;
uint16_t timeout_ms; uint32_t sum_rtt; /* в 0.1 ms */
uint16_t interval_ms; uint16_t timeout_ms;
void* next_timer;
}; };
// ==================================== блок обработки запросов на удаленные пинги
struct route_ping_pending { struct route_ping_pending {
struct route_ping_pending* next; struct route_ping_pending* next;
struct ROUTE_BGP* bgp; struct ROUTE_BGP* bgp;
@ -62,10 +43,7 @@ struct route_ping_pending {
uint8_t cancelled; uint8_t cancelled;
}; };
static void route_ping_send_resp(struct route_ping_req* req, uint16_t avg_rtt); // если удаленный узел долго не отвечает - вызываем коллбэк по таймауту
static void route_ping_next(void* arg);
static void route_ping_finish(struct route_ping_req* req);
static void route_ping_pending_timeout(void* arg) { static void route_ping_pending_timeout(void* arg) {
struct route_ping_pending* p = (struct route_ping_pending*)arg; struct route_ping_pending* p = (struct route_ping_pending*)arg;
if (!p) return; if (!p) return;
@ -89,142 +67,7 @@ static void route_ping_pending_timeout(void* arg) {
u_free(p); u_free(p);
} }
static void route_ping_series_free(struct route_ping_series* series) { // отправить запрос удаленному узлу "пропингуй такой-то узел"
if (!series) return;
if (series->next_timer) {
uasync_cancel_timeout(series->req->reply_conn->instance->ua, series->next_timer);
series->next_timer = NULL;
}
u_free(series);
}
static void route_ping_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 route_ping_series* series = (struct route_ping_series*)arg;
if (!series) return;
series->count_sent++;
if (success) {
series->count_ok++;
series->rtt_sum += rtt;
}
if (series->count_sent < series->count_total && series->req->reply_conn) {
if (success) {
// Ответ получен — следующий пинг сразу (burst mode)
route_ping_next(series);
} else {
// Таймаут — ждем interval_ms перед retry
struct UTUN_INSTANCE* inst = series->req->reply_conn->instance;
series->next_timer = uasync_set_timeout(inst->ua, series->interval_ms * 10, series, route_ping_next);
}
return;
}
// серия завершена
if (series->count_ok > 0) {
series->sock_ctx->avg_rtt = (uint16_t)(series->rtt_sum / series->count_ok);
} else {
series->sock_ctx->avg_rtt = series->timeout_ms * 10; // максимум при полном провале
}
series->sock_ctx->count_sent = series->count_sent;
series->sock_ctx->count_ok = series->count_ok;
series->req->completed_count++;
if (series->req->completed_count >= series->req->socket_count) {
route_ping_finish(series->req);
}
u_free(series);
}
static void route_ping_next(void* arg) {
struct route_ping_series* series = (struct route_ping_series*)arg;
if (!series) return;
series->next_timer = NULL;
if (!series->req->reply_conn) {
u_free(series);
return;
}
int ret = etcp_send_ping_to_socket(series->req->reply_conn->instance, series->local_sock,
series->pubkey, &series->target_addr,
series->timeout_ms, route_ping_cb, series, NULL, 0);
if (ret != 0) {
// не удалось отправить, считаем как fail и продолжаем/завершаем
series->count_sent++;
if (series->count_sent < series->count_total) {
struct UTUN_INSTANCE* inst = series->req->reply_conn->instance;
series->next_timer = uasync_set_timeout(inst->ua, series->interval_ms * 10, series, route_ping_next);
} else {
series->sock_ctx->avg_rtt = series->timeout_ms * 10;
series->sock_ctx->count_sent = series->count_sent;
series->sock_ctx->count_ok = series->count_ok;
series->req->completed_count++;
if (series->req->completed_count >= series->req->socket_count) {
route_ping_finish(series->req);
}
u_free(series);
}
}
}
static void route_ping_finish(struct route_ping_req* req) {
if (!req) return;
uint32_t best_metric = 0xFFFFFFFF;
uint16_t best_avg_rtt = 0;
uint8_t best_count_sent = 0;
uint8_t best_count_ok = 0;
struct ETCP_SOCKET* best_sock = NULL;
for (uint8_t i = 0; i < req->socket_count; i++) {
struct route_ping_sock_ctx* sc = &req->sockets[i];
uint8_t loss = (sc->count_sent > 0) ? ((sc->count_sent - sc->count_ok) * 100 / sc->count_sent) : 0;
uint32_t metric = (uint32_t)sc->avg_rtt * (loss * 3 + 1);
if (metric < best_metric) {
best_metric = metric;
best_avg_rtt = sc->avg_rtt;
best_count_sent = sc->count_sent;
best_count_ok = sc->count_ok;
best_sock = sc->local_sock;
}
}
if (req->target_node) {
req->target_node->last_ping_time = get_time_tb();
req->target_node->last_rtt = best_avg_rtt;
req->target_node->best_socket = best_sock;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "node=%016llx best_rtt=%u best_sock=%p metric=%u",
(unsigned long long)req->target_node->node.node_id, (unsigned)best_avg_rtt,
(void*)best_sock, (unsigned)best_metric);
}
route_ping_send_resp(req, best_avg_rtt);
u_free(req);
}
static void route_ping_send_resp(struct route_ping_req* req, uint16_t avg_rtt) {
if (!req || !req->reply_conn) return;
struct BGP_PING_RESPONSE* resp = u_calloc(1, sizeof(struct BGP_PING_RESPONSE));
if (!resp) return;
resp->cmd = ETCP_ID_ROUTE_ENTRY;
resp->subcmd = ROUTE_SUBCMD_PING_RESP;
resp->request_id = req->request_id;
// суммируем по всем сокетам для отчёта
uint8_t total_sent = 0, total_ok = 0;
for (uint8_t i = 0; i < req->socket_count; i++) {
total_sent += req->sockets[i].count_sent;
total_ok += req->sockets[i].count_ok;
}
resp->count_sent = total_sent;
resp->count_ok = total_ok;
resp->avg_rtt = avg_rtt;
memcpy(resp->recv_ipv4, req->recv_ipv4, 4);
resp->recv_port = htons(req->recv_port);
struct ll_entry* e = queue_entry_new(0);
if (!e) {
u_free(resp);
return;
}
e->dgram = (uint8_t*)resp;
e->len = sizeof(struct BGP_PING_RESPONSE);
etcp_send(req->reply_conn, e);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "request_id=%016llx sent=%u ok=%u avg_rtt=%u",
(unsigned long long)req->request_id, (unsigned)total_sent, (unsigned)total_ok, (unsigned)avg_rtt);
}
int route_ping_send_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, uint64_t node_id, int route_ping_send_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, uint64_t node_id,
uint8_t count, uint16_t interval_ms, uint16_t timeout_ms, uint8_t count, uint16_t interval_ms, uint16_t timeout_ms,
uint16_t wait_timeout_ms, route_ping_callback_t cb, void* arg) { uint16_t wait_timeout_ms, route_ping_callback_t cb, void* arg) {
@ -282,6 +125,44 @@ int route_ping_send_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, uint64
return 0; return 0;
} }
// прошел ответ "серия пигнов на удаленном узле завершена"
void route_ping_handle_resp(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len) {
if (!bgp || !from_conn || !data || len < sizeof(struct BGP_PING_RESPONSE)) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args");
return;
}
const struct BGP_PING_RESPONSE* resp = (const struct BGP_PING_RESPONSE*)data;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "from=%s request_id=%016llx sent=%u ok=%u avg_rtt=%u",
from_conn->log_name, (unsigned long long)resp->request_id,
(unsigned)resp->count_sent, (unsigned)resp->count_ok, (unsigned)resp->avg_rtt);
struct route_ping_pending** cur = &bgp->ping_pending;
while (*cur) {
struct route_ping_pending* p = *cur;
if (p->request_id == resp->request_id) {
*cur = p->next;
if (p->timeout_timer) {
uasync_cancel_timeout(bgp->instance->ua, p->timeout_timer);
p->timeout_timer = NULL;
}
if (p->callback) {
int success = (resp->count_ok > 0) ? 1 : 0;
uint32_t recv_ip = (resp->recv_ipv4[0] << 24) | (resp->recv_ipv4[1] << 16) |
(resp->recv_ipv4[2] << 8) | resp->recv_ipv4[3];
uint16_t recv_port = ntohs(resp->recv_port);
p->callback(success, resp->avg_rtt, resp->count_sent, resp->count_ok, recv_ip, recv_port, p->arg);
}
u_free(p);
return;
}
cur = &(*cur)->next;
}
DEBUG_WARN(DEBUG_CATEGORY_BGP, "request_id=%016llx not found in pending list",
(unsigned long long)resp->request_id);
}
// ========================================================================
int route_ping_send_req_addr(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, uint64_t node_id, int route_ping_send_req_addr(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, uint64_t node_id,
uint32_t target_ip, uint16_t target_port, uint32_t target_ip, uint16_t target_port,
uint8_t count, uint16_t interval_ms, uint16_t timeout_ms, uint8_t count, uint16_t interval_ms, uint16_t timeout_ms,
@ -347,263 +228,159 @@ int route_ping_send_req_addr(struct ROUTE_BGP* bgp, struct ETCP_CONN* to_conn, u
return 0; return 0;
} }
void route_ping_destroy_pending(struct ROUTE_BGP* bgp) {
if (!bgp) return; static void route_ping_series_finish(struct route_ping_series_ctx* ctx) {
while (bgp->ping_pending) { if (!ctx || !ctx->reply_conn) {
struct route_ping_pending* p = bgp->ping_pending; u_free(ctx);
bgp->ping_pending = p->next; return;
p->next = NULL; }
if (p->timeout_timer) {
err_t rc = uasync_cancel_timeout(bgp->instance->ua, p->timeout_timer); uint16_t avg_rtt = (ctx->count_ok > 0)
p->timeout_timer = NULL; ? (uint16_t)(ctx->sum_rtt / ctx->count_ok)
if (rc == ERR_OK) { : (uint16_t)(ctx->timeout_ms * 10); /* таймаут как "плохо" */
u_free(p);
} else { /* Отправляем ответ по BGP */
p->cancelled = 1; struct BGP_PING_RESPONSE* resp = u_calloc(1, sizeof(struct BGP_PING_RESPONSE));
} if (resp) {
resp->cmd = ETCP_ID_ROUTE_ENTRY;
resp->subcmd = ROUTE_SUBCMD_PING_RESP;
resp->request_id = ctx->request_id;
resp->count_sent = ctx->count_sent;
resp->count_ok = ctx->count_ok;
resp->avg_rtt = avg_rtt;
/* recv_ipv4/port = 0 (не используется) */
memset(resp->recv_ipv4, 0, 4);
resp->recv_port = 0;
struct ll_entry* e = queue_entry_new(0);
if (e) {
e->dgram = (uint8_t*)resp;
e->len = sizeof(struct BGP_PING_RESPONSE);
etcp_send(ctx->reply_conn, e);
DEBUG_INFO(DEBUG_CATEGORY_BGP,
"PING series done request_id=%016llx sent=%u ok=%u avg_rtt=%u",
(unsigned long long)ctx->request_id,
(unsigned)ctx->count_sent,
(unsigned)ctx->count_ok,
(unsigned)avg_rtt);
} else { } else {
u_free(p); u_free(resp);
} }
} }
u_free(ctx);
} }
void route_ping_handle_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len) {// получен запрос сделать пинг. /* Callback одного пинга из серии */
const size_t base_len = offsetof(struct BGP_PING_REQUEST, target_ipv4); static void route_ping_single_cb(int success,
const size_t addr_len = offsetof(struct BGP_PING_REQUEST, pubkey); uint16_t rtt,
if (!bgp || !from_conn || !data || (len != base_len && len != addr_len && len != sizeof(struct BGP_PING_REQUEST))) { void* arg,
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args len=%zu", len); uint64_t nonce,
return; const uint8_t* resp_data,
} size_t resp_data_len) {
const struct BGP_PING_REQUEST* req_pkt = (const struct BGP_PING_REQUEST*)data; (void)nonce; (void)resp_data; (void)resp_data_len;
// Check for custom target IP:port first (for NAT detection with STUN) struct route_ping_series_ctx* ctx = (struct route_ping_series_ctx*)arg;
struct NODEINFO_IPV4_SOCKET custom_socket; if (!ctx) return;
const struct NODEINFO_IPV4_SOCKET* sockets_v4 = NULL;
const struct NODEINFO_IPV6_SOCKET* sockets_v6 = NULL; /* Обновляем статистику */
int sock_count_v4 = 0; ctx->count_sent++;
int sock_count_v6 = 0; if (success) {
struct NODEINFO_Q* target = NULL; ctx->count_ok++;
ctx->sum_rtt += rtt;
const uint8_t* embedded_pubkey = NULL; }
if (len == sizeof(struct BGP_PING_REQUEST) || len == addr_len) {
// Extended packet with custom target (with or without pubkey) /* Если ещё не все пакеты отправлены — сразу шлём следующий */
custom_socket.addr[0] = req_pkt->target_ipv4[0]; if (ctx->count_sent < ctx->count_total) {
custom_socket.addr[1] = req_pkt->target_ipv4[1]; int ret = etcp_send_ping_to_socket(
custom_socket.addr[2] = req_pkt->target_ipv4[2]; ctx->reply_conn->instance,
custom_socket.addr[3] = req_pkt->target_ipv4[3]; ctx->local_sock,
custom_socket.port = req_pkt->target_port; ctx->pubkey,
custom_socket.type = 0; &ctx->target_addr,
custom_socket.id = 0; ctx->timeout_ms,
sockets_v4 = &custom_socket; route_ping_single_cb,
sock_count_v4 = 1; ctx,
DEBUG_INFO(DEBUG_CATEGORY_BGP, "using custom target %d.%d.%d.%d:%u for node %016llx", NULL, 0);
req_pkt->target_ipv4[0], req_pkt->target_ipv4[1], req_pkt->target_ipv4[2], req_pkt->target_ipv4[3],
htons(req_pkt->target_port), (unsigned long long)req_pkt->node_id); if (ret != 0) {
// Check if embedded pubkey is non-zero (only when full packet length) /* Не смогли отправить следующий — завершаем серию досрочно */
if (len == sizeof(struct BGP_PING_REQUEST)) { route_ping_series_finish(ctx);
int pubkey_zero = 1;
for (int i = 0; i < SC_PUBKEY_SIZE; i++) {
if (req_pkt->pubkey[i] != 0) { pubkey_zero = 0; break; }
}
if (!pubkey_zero) {
embedded_pubkey = req_pkt->pubkey;
}
}
// For custom target, we still need target node for pubkey lookup if not embedded
target = route_bgp_get_node(bgp, req_pkt->node_id);
} else if (len == base_len) {
// Base packet without custom target - use nodeinfo sockets
target = route_bgp_get_node(bgp, req_pkt->node_id);
if (!target) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "node %016llx not found",
(unsigned long long)req_pkt->node_id);
// отправляем ответ с нулями
struct route_ping_req* req = u_calloc(1, sizeof(struct route_ping_req));
if (req) {
req->reply_conn = from_conn;
req->request_id = req_pkt->request_id;
req->socket_count = 0;
route_ping_send_resp(req, 0);
u_free(req);
}
return;
} }
sock_count_v4 = get_node_v4_sockets(target, &sockets_v4); /* else: следующий пинг запущен, ждём его callback */
sock_count_v6 = get_node_v6_sockets(target, &sockets_v6);
} else { } else {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "unexpected packet size %zu", len); /* Все пакеты обработаны */
return; route_ping_series_finish(ctx);
} }
}
if ((sock_count_v4 <= 0 || !sockets_v4) && (sock_count_v6 <= 0 || !sockets_v6)) { // Обработчик PING_REQ (нам запросили "пропингуй узел и верни результат")
DEBUG_WARN(DEBUG_CATEGORY_BGP, "no sockets for node %016llx", void route_ping_handle_req(struct ROUTE_BGP* bgp,
(unsigned long long)req_pkt->node_id); struct ETCP_CONN* from_conn,
struct route_ping_req* req = u_calloc(1, sizeof(struct route_ping_req)); const uint8_t* data,
if (req) { size_t len) {
req->reply_conn = from_conn; if (!bgp || !from_conn || !data || len < sizeof(struct BGP_PING_REQUEST)) {
req->request_id = req_pkt->request_id; DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping req: bad args len=%zu", len);
req->socket_count = 0;
route_ping_send_resp(req, 0);
u_free(req);
}
return; return;
} }
// считаем подходящие локальные сокеты (v4 + v6) const struct BGP_PING_REQUEST* req_pkt = (const struct BGP_PING_REQUEST*)data;
uint8_t local_v4_count = 0;
uint8_t local_v6_count = 0; /* Создаём контекст серии */
struct ETCP_SOCKET* ls = bgp->instance->etcp_sockets; struct route_ping_series_ctx* ctx = u_calloc(1, sizeof(struct route_ping_series_ctx));
while (ls) { if (!ctx) {
if (ls->local_addr.ss_family == AF_INET) local_v4_count++; DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping req: alloc ctx failed");
else if (ls->local_addr.ss_family == AF_INET6) local_v6_count++;
ls = ls->next;
}
uint8_t local_count = local_v4_count + local_v6_count;
if (local_count == 0) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "no local sockets");
struct route_ping_req* req = u_calloc(1, sizeof(struct route_ping_req));
if (req) {
req->reply_conn = from_conn;
req->request_id = req_pkt->request_id;
req->socket_count = 0;
route_ping_send_resp(req, 0);
u_free(req);
}
return;
}
struct route_ping_req* req = u_calloc(1, sizeof(struct route_ping_req) + local_count * sizeof(struct route_ping_sock_ctx));
if (!req) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "alloc failed");
return; return;
} }
req->reply_conn = from_conn;
req->request_id = req_pkt->request_id;
req->target_node = target;
req->timeout_ms = req_pkt->timeout_ms;
req->socket_count = local_count;
// Сохраняем IP:port запрашивающего (из первого линка) для STUN-ответа
struct ETCP_LINK* l = from_conn->links;
if (l && l->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&l->remote_addr;
req->recv_ipv4[0] = (ntohl(sin->sin_addr.s_addr) >> 24) & 0xFF;
req->recv_ipv4[1] = (ntohl(sin->sin_addr.s_addr) >> 16) & 0xFF;
req->recv_ipv4[2] = (ntohl(sin->sin_addr.s_addr) >> 8) & 0xFF;
req->recv_ipv4[3] = ntohl(sin->sin_addr.s_addr) & 0xFF;
req->recv_port = ntohs(sin->sin_port);
} else {
memset(req->recv_ipv4, 0, 4);
req->recv_port = 0;
}
ls = bgp->instance->etcp_sockets; ctx->reply_conn = from_conn;
uint8_t idx = 0; ctx->request_id = req_pkt->request_id;
while (ls) { ctx->count_total = req_pkt->count;
struct sockaddr_storage target_addr; ctx->timeout_ms = req_pkt->timeout_ms;
memset(&target_addr, 0, sizeof(target_addr)); memcpy(ctx->pubkey, req_pkt->pubkey, SC_PUBKEY_SIZE);
if (ls->local_addr.ss_family == AF_INET) { /* Целевой адрес */
if (sock_count_v4 > 0 && sockets_v4) { struct sockaddr_in* sin = (struct sockaddr_in*)&ctx->target_addr;
struct sockaddr_in* sin = (struct sockaddr_in*)&target_addr; sin->sin_family = AF_INET;
sin->sin_family = AF_INET; memcpy(&sin->sin_addr.s_addr, req_pkt->target_ipv4, 4);
memcpy(&sin->sin_addr, sockets_v4[0].addr, 4); sin->sin_port = req_pkt->target_port;
sin->sin_port = htons(sockets_v4[0].port);
} else {
ls = ls->next;
continue;
}
} else if (ls->local_addr.ss_family == AF_INET6) {
if (sock_count_v6 > 0 && sockets_v6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&target_addr;
sin6->sin6_family = AF_INET6;
memcpy(&sin6->sin6_addr, sockets_v6[0].addr, 16);
sin6->sin6_port = htons(sockets_v6[0].port);
} else {
ls = ls->next;
continue;
}
} else {
ls = ls->next;
continue;
}
req->sockets[idx].local_sock = ls; /* Ищем первый IPv4-сокет (как в старом коде) */
struct route_ping_series* series = u_calloc(1, sizeof(struct route_ping_series)); struct ETCP_SOCKET* ls = bgp->instance->etcp_sockets;
if (series) { while (ls) {
series->req = req; if (ls->local_addr.ss_family == AF_INET) {
series->sock_ctx = &req->sockets[idx]; ctx->local_sock = ls;
series->local_sock = ls; break;
series->target_addr = target_addr;
const uint8_t* pubkey_to_use = target ? target->node.public_key : embedded_pubkey;
if (!pubkey_to_use) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "no pubkey available for node %016llx",
(unsigned long long)req_pkt->node_id);
series->sock_ctx->avg_rtt = req_pkt->timeout_ms * 10;
series->sock_ctx->count_sent = req_pkt->count;
req->completed_count++;
u_free(series);
idx++;
ls = ls->next;
continue;
}
memcpy(series->pubkey, pubkey_to_use, SC_PUBKEY_SIZE);
series->count_total = req_pkt->count;
series->count_sent = 0;
series->count_ok = 0;
series->rtt_sum = 0;
series->timeout_ms = req_pkt->timeout_ms;
series->interval_ms = req_pkt->interval_ms;
int ret = etcp_send_ping_to_socket(bgp->instance, ls, pubkey_to_use,
&target_addr, req_pkt->timeout_ms, route_ping_cb, series, NULL, 0);
if (ret != 0) {
series->sock_ctx->avg_rtt = req_pkt->timeout_ms * 10;
series->sock_ctx->count_sent = req_pkt->count;
req->completed_count++;
u_free(series);
}
} else {
req->sockets[idx].avg_rtt = req_pkt->timeout_ms * 10;
req->sockets[idx].count_sent = req_pkt->count;
req->completed_count++;
} }
idx++;
ls = ls->next; ls = ls->next;
} }
if (req->completed_count >= req->socket_count) {
route_ping_finish(req);
}
}
void route_ping_handle_resp(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len) { if (!ctx->local_sock) {
if (!bgp || !from_conn || !data || len < sizeof(struct BGP_PING_RESPONSE)) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "ping req: no IPv4 socket");
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); route_ping_series_finish(ctx); /* отправит 0/0 */
return; return;
} }
const struct BGP_PING_RESPONSE* resp = (const struct BGP_PING_RESPONSE*)data;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "from=%s request_id=%016llx sent=%u ok=%u avg_rtt=%u",
from_conn->log_name, (unsigned long long)resp->request_id,
(unsigned)resp->count_sent, (unsigned)resp->count_ok, (unsigned)resp->avg_rtt);
struct route_ping_pending** cur = &bgp->ping_pending; DEBUG_INFO(DEBUG_CATEGORY_BGP,
while (*cur) { "PING series start request_id=%016llx target=%s:%u count=%u timeout=%u",
struct route_ping_pending* p = *cur; (unsigned long long)ctx->request_id,
if (p->request_id == resp->request_id) { ip_to_str(&sin->sin_addr, AF_INET).str,
*cur = p->next; ntohs(sin->sin_port),
if (p->timeout_timer) { (unsigned)ctx->count_total,
uasync_cancel_timeout(bgp->instance->ua, p->timeout_timer); (unsigned)ctx->timeout_ms);
p->timeout_timer = NULL;
} /* Запускаем первый пинг (дальше цепочка через callback) */
if (p->callback) { int ret = etcp_send_ping_to_socket(
int success = (resp->count_ok > 0) ? 1 : 0; bgp->instance,
uint32_t recv_ip = (resp->recv_ipv4[0] << 24) | (resp->recv_ipv4[1] << 16) | ctx->local_sock,
(resp->recv_ipv4[2] << 8) | resp->recv_ipv4[3]; ctx->pubkey,
uint16_t recv_port = ntohs(resp->recv_port); &ctx->target_addr,
p->callback(success, resp->avg_rtt, resp->count_sent, resp->count_ok, recv_ip, recv_port, p->arg); ctx->timeout_ms,
} route_ping_single_cb,
u_free(p); ctx,
return; NULL, 0);
}
cur = &(*cur)->next; if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping req: cannot start first ping");
route_ping_series_finish(ctx);
} }
DEBUG_WARN(DEBUG_CATEGORY_BGP, "request_id=%016llx not found in pending list",
(unsigned long long)resp->request_id);
} }

1
src/route_ping.h

@ -14,6 +14,7 @@ struct BGP_PING_REQUEST {
uint64_t request_id; // для корреляции uint64_t request_id; // для корреляции
uint64_t node_id; // целевой узел uint64_t node_id; // целевой узел
uint8_t count; // число пингов uint8_t count; // число пингов
uint8_t socket_id; // id сокета пингуемого узла
uint16_t interval_ms; // интервал между пингами uint16_t interval_ms; // интервал между пингами
uint16_t timeout_ms; // таймаут одного пинга uint16_t timeout_ms; // таймаут одного пинга
uint8_t target_ipv4[4]; // custom target IP (0 = use nodeinfo sockets) uint8_t target_ipv4[4]; // custom target IP (0 = use nodeinfo sockets)

2
tools/etcpmon/etcpmon_protocol.h

@ -33,7 +33,7 @@ extern "C" {
#define NAT_VERIFIED_DIRECT 7 #define NAT_VERIFIED_DIRECT 7
#define ETCPMON_MAX_MSG_SIZE 4096 #define ETCPMON_MAX_MSG_SIZE 4096
#define ETCPMON_MAX_CONN_NAME 32 #define ETCPMON_MAX_CONN_NAME 32
#define ETCPMON_MAX_CONNECTIONS 256 #define ETCPMON_MAX_CONNECTIONS 250
#define ETCPMON_MAX_LINKS 16 #define ETCPMON_MAX_LINKS 16
/* Update interval in milliseconds */ /* Update interval in milliseconds */

Loading…
Cancel
Save