Browse Source

add ping

congestion
Evgeny 6 months ago
parent
commit
ffbc0da86f
  1. 4
      .gitignore
  2. 1
      src/Makefile.am
  3. 106
      src/etcp_connections.c
  4. 13
      src/etcp_connections.h
  5. 9
      src/route_bgp.c
  6. 4
      src/route_bgp.h
  7. 78
      src/route_node.c
  8. 13
      src/route_node.h
  9. 421
      src/route_ping.c
  10. 45
      src/route_ping.h
  11. 6
      tests/Makefile.am
  12. 10
      tests/test_etcp_ping.c

4
.gitignore vendored

@ -9,7 +9,6 @@ utun
# binary files
*.exe
**/*[^.]
# Autotools
Makefile
@ -57,7 +56,8 @@ tests/*.trs
tests/logs/
# Test binaries (files without extension in tests/)
tests/test_*
tests/bench_*
# All log files
*.log

1
src/Makefile.am

@ -8,6 +8,7 @@ utun_CORE_SOURCES = \
config_updater.c \
route_lib.c \
route_bgp.c \
route_ping.c \
route_node.c \
routing.c \
tun_if.c \

106
src/etcp_connections.c

@ -908,6 +908,13 @@ static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, socket_t fd, sc_context_
return (int)sent;
}
static size_t etcp_build_ping_response_data(const uint8_t* req_data, size_t req_data_len,
uint8_t* resp_buf, size_t resp_buf_len) {
if (req_data_len > resp_buf_len) req_data_len = resp_buf_len;
if (req_data_len && req_data) memcpy(resp_buf, req_data, req_data_len);
return req_data_len;
}
static void ping_timeout_cbk(void* arg) {
struct PING_CONTEXT* ctx = (struct PING_CONTEXT*)arg;
if (!ctx || !ctx->cb) return;
@ -915,7 +922,7 @@ static void ping_timeout_cbk(void* arg) {
uasync_cancel_timeout(ctx->instance->ua, ctx->timeout_timer);
ctx->timeout_timer = NULL;
}
ctx->cb(0, ctx->arg, ctx->nonce);
ctx->cb(0, 0, ctx->arg, ctx->nonce, NULL, 0);
if (ctx->instance->pending_pings == ctx) {
ctx->instance->pending_pings = ctx->next;
} else {
@ -923,18 +930,20 @@ static void ping_timeout_cbk(void* arg) {
while (prev && prev->next != ctx) prev = prev->next;
if (prev) prev->next = ctx->next;
}
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
}
int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bin, const struct sockaddr_storage* addr, int timeout_ms, etcp_ping_callback_t cb, void* user_arg) {
if (!instance || !peer_pubkey_bin || !addr || timeout_ms <= 0 || !cb) {
int etcp_send_ping_to_socket(struct UTUN_INSTANCE* instance, struct ETCP_SOCKET* e_sock,
const uint8_t* peer_pubkey_bin, const struct sockaddr_storage* addr,
int timeout_ms, etcp_ping_callback_t cb, void* user_arg,
const uint8_t* user_data, size_t user_data_len) {
if (!instance || !e_sock || !peer_pubkey_bin || !addr || timeout_ms <= 0 || !cb) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "bad args");
return -1;
}
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock && e_sock->local_addr.ss_family != addr->ss_family) e_sock = e_sock->next;
if (!e_sock) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "no socket");
if (user_data_len > PACKET_DATA_SIZE - 24 - 2 - SC_PUBKEY_ENC_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "user_data too long");
return -2;
}
struct PING_CONTEXT* ctx = u_malloc(sizeof(struct PING_CONTEXT));
@ -948,32 +957,54 @@ int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bi
ctx->arg = user_arg;
ctx->nonce = get_current_timestamp() ^ (uint64_t)rand();
ctx->timeout_timer = NULL;
ctx->send_time = get_time_tb();
ctx->user_data = NULL;
ctx->user_data_len = 0;
if (user_data_len > 0) {
ctx->user_data = u_malloc(user_data_len);
if (!ctx->user_data) {
u_free(ctx);
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "malloc user_data");
return -3;
}
memcpy(ctx->user_data, user_data, user_data_len);
ctx->user_data_len = user_data_len;
}
struct ETCP_DGRAM* dgram = u_malloc(PACKET_DATA_SIZE);
if (!dgram) {
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "malloc dgram");
return -4;
}
dgram->link = NULL;
dgram->noencrypt_len = SC_PUBKEY_ENC_SIZE;
dgram->data_len = 24;
size_t offset = 0;
uint8_t* p = dgram->data;
*p++ = ETCP_PING;
uint64_t nid = htobe64(instance->node_id);
memcpy(p, &nid, 8); p += 8;
uint64_t nonce_be = htobe64(ctx->nonce);
memcpy(p, &nonce_be, 8);
memcpy(p, &nonce_be, 8); p += 8;
uint16_t ulen_be = htobe16((uint16_t)user_data_len);
memcpy(p, &ulen_be, 2); p += 2;
if (user_data_len) {
memcpy(p, user_data, user_data_len);
p += user_data_len;
}
uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE];
random_bytes(salt, sizeof(salt));
memcpy(dgram->data + 24, salt, SC_PUBKEY_ENC_SALT_SIZE);
memcpy(p, salt, SC_PUBKEY_ENC_SALT_SIZE); p += SC_PUBKEY_ENC_SALT_SIZE;
uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE];
sc_obfuscate_pubkey(salt, peer_pubkey_bin, instance->my_keys.public_key, obfuscated_pubkey);
memcpy(dgram->data + 24 + SC_PUBKEY_ENC_SALT_SIZE, obfuscated_pubkey, SC_PUBKEY_SIZE);
dgram->data_len = 24 + SC_PUBKEY_ENC_SIZE;
memcpy(p, obfuscated_pubkey, SC_PUBKEY_SIZE); p += SC_PUBKEY_SIZE;
dgram->data_len = (uint16_t)(p - dgram->data);
struct secure_channel sc;
sc_init_ctx(&sc, &instance->my_keys);
if (sc_set_peer_public_key(&sc, peer_pubkey_bin, SC_PEER_PUBKEY_BIN) != SC_OK) {
u_free(dgram); u_free(ctx);
u_free(dgram);
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "set key failed");
return -5;
}
@ -987,10 +1018,29 @@ int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bi
last->next = ctx;
}
ctx->timeout_timer = uasync_set_timeout(instance->ua, timeout_ms * 10, ctx, ping_timeout_cbk);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "ping sent nonce=%016llx timeout=%d", (unsigned long long)ctx->nonce, timeout_ms);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "ping sent nonce=%016llx timeout=%d ulen=%zu",
(unsigned long long)ctx->nonce, timeout_ms, user_data_len);
return 0;
}
int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bin,
const struct sockaddr_storage* addr, int timeout_ms,
etcp_ping_callback_t cb, void* user_arg,
const uint8_t* user_data, size_t user_data_len) {
if (!instance || !peer_pubkey_bin || !addr || timeout_ms <= 0 || !cb) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "bad args");
return -1;
}
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock && e_sock->local_addr.ss_family != addr->ss_family) e_sock = e_sock->next;
if (!e_sock) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "no socket");
return -2;
}
return etcp_send_ping_to_socket(instance, e_sock, peer_pubkey_bin, addr, timeout_ms,
cb, user_arg, user_data, user_data_len);
}
void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connections_read_callback fd=%d, socket=%p", fd, arg);
@ -1098,24 +1148,29 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
uint64_t peer_id = be64toh(*(uint64_t*)ack_hdr->id);
if (ack_hdr->code == ETCP_PING) {
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = be16toh(*(uint16_t*)(pkt->data + 17));
const uint8_t* udata = (ulen > 0) ? (pkt->data + 19) : NULL;
struct ETCP_DGRAM* resp = u_malloc(PACKET_DATA_SIZE);
if (resp) {
resp->link = NULL;
resp->noencrypt_len = SC_PUBKEY_ENC_SIZE;
resp->data_len = 24;
uint8_t* p = resp->data;
*p++ = ETCP_PONG;
uint64_t nid = htobe64(e_sock->instance->node_id);
memcpy(p, &nid, 8); p += 8;
uint64_t n = htobe64(nonce);
memcpy(p, &n, 8);
memcpy(p, &n, 8); p += 8;
uint16_t resp_ulen_be = htobe16(ulen);
memcpy(p, &resp_ulen_be, 2); p += 2;
size_t copied = etcp_build_ping_response_data(udata, ulen, p, PACKET_DATA_SIZE - (p - resp->data) - SC_PUBKEY_ENC_SIZE);
p += copied;
uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE];
random_bytes(salt, sizeof(salt));
memcpy(resp->data + 24, salt, SC_PUBKEY_ENC_SALT_SIZE);
memcpy(p, salt, SC_PUBKEY_ENC_SALT_SIZE); p += SC_PUBKEY_ENC_SALT_SIZE;
uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE];
sc_obfuscate_pubkey(salt, decrypted_pubkey, e_sock->instance->my_keys.public_key, obfuscated_pubkey);
memcpy(resp->data + 24 + SC_PUBKEY_ENC_SALT_SIZE, obfuscated_pubkey, SC_PUBKEY_SIZE);
resp->data_len = 24 + SC_PUBKEY_ENC_SIZE;
memcpy(p, obfuscated_pubkey, SC_PUBKEY_SIZE); p += SC_PUBKEY_SIZE;
resp->data_len = (uint16_t)(p - resp->data);
sc_context_t resp_sc;
sc_init_ctx(&resp_sc, &e_sock->instance->my_keys);
if (sc_set_peer_public_key(&resp_sc, decrypted_pubkey, SC_PEER_PUBKEY_BIN) == SC_OK) {
@ -1128,6 +1183,14 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
}
if (ack_hdr->code == ETCP_PONG) {
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = 0;
const uint8_t* udata = NULL;
if (pkt->data_len >= 19) {
ulen = be16toh(*(uint16_t*)(pkt->data + 17));
if (ulen > 0 && pkt->data_len >= 19 + ulen) {
udata = pkt->data + 19;
}
}
struct PING_CONTEXT* ctx = e_sock->instance->pending_pings;
struct PING_CONTEXT* prev = NULL;
while (ctx) {
@ -1138,7 +1201,10 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
uasync_cancel_timeout(e_sock->instance->ua, ctx->timeout_timer);
ctx->timeout_timer = NULL;
}
ctx->cb(1, ctx->arg, nonce);
uint64_t now = get_time_tb();
uint16_t rtt = (now >= ctx->send_time) ? (uint16_t)(now - ctx->send_time) : 0;
ctx->cb(1, rtt, ctx->arg, nonce, udata, ulen);
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
break;
}

13
src/etcp_connections.h

@ -35,7 +35,8 @@ struct ETCP_DGRAM {// пакет (незашифрованный)
};
#pragma pack(pop)
typedef void (*etcp_ping_callback_t)(int success, void* arg, uint64_t nonce);
typedef void (*etcp_ping_callback_t)(int success, uint16_t rtt, void* arg, uint64_t nonce,
const uint8_t* resp_data, size_t resp_data_len);
struct PING_CONTEXT {
struct PING_CONTEXT* next;
struct UTUN_INSTANCE* instance;
@ -43,6 +44,9 @@ struct PING_CONTEXT {
void* arg;
uint64_t nonce;
void* timeout_timer;
uint64_t send_time; // время отправки пинга в 0.1ms
uint8_t* user_data;
size_t user_data_len;
};
// список активных подключений которые обслуживает сокет. каждый сокет может обслуживать много подключений
@ -205,7 +209,12 @@ void start_stats_timer(struct ETCP_LINK* link);
int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bin,
const struct sockaddr_storage* addr, int timeout_ms,
etcp_ping_callback_t cb, void* user_arg);
etcp_ping_callback_t cb, void* user_arg,
const uint8_t* user_data, size_t user_data_len);
int etcp_send_ping_to_socket(struct UTUN_INSTANCE* instance, struct ETCP_SOCKET* e_sock,
const uint8_t* peer_pubkey_bin, const struct sockaddr_storage* addr,
int timeout_ms, etcp_ping_callback_t cb, void* user_arg,
const uint8_t* user_data, size_t user_data_len);
void etcp_connections_read_callback_socket(socket_t sock, void* arg);

9
src/route_bgp.c

@ -17,6 +17,7 @@
#include "route_node.h"
#include "route_lib.h"
#include "route_bgp.h"
#include "route_ping.h"
// ============================================================================
@ -150,6 +151,10 @@ static void route_bgp_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
route_bgp_process_withdraw(bgp, from_conn, data, entry->len);
} else if (subcmd == ROUTE_SUBCMD_REQUEST_TABLE) {
route_bgp_handle_request_table(bgp, from_conn);
} else if (subcmd == ROUTE_SUBCMD_PING_REQ) {
route_ping_handle_req(bgp, from_conn, data, entry->len);
} else if (subcmd == ROUTE_SUBCMD_PING_RESP) {
route_ping_handle_resp(bgp, from_conn, data, entry->len);
}
queue_dgram_free(entry);
@ -207,6 +212,8 @@ struct ROUTE_BGP* route_bgp_init(struct UTUN_INSTANCE* instance) {
}
bgp->instance = instance;
bgp->ping_pending = NULL;
bgp->next_ping_req_id = 1;
// Create nodes queue with hash support for fast lookup by node_id
bgp->nodes = queue_new(instance->ua, BGP_NODES_HASH_SIZE, "bgp_nodes");
@ -270,6 +277,8 @@ void route_bgp_destroy(struct UTUN_INSTANCE* instance) {
etcp_unbind(instance, ETCP_ID_ROUTE_ENTRY);
route_ping_destroy_pending(instance->bgp);
struct ll_entry* e;
while ((e = queue_data_get(instance->bgp->senders_list)) != NULL) {
queue_entry_free(e);

4
src/route_bgp.h

@ -52,11 +52,15 @@ struct ROUTE_BGP_CONN_ITEM {
struct ETCP_CONN* conn;
};
struct route_ping_pending;
struct ROUTE_BGP {
struct UTUN_INSTANCE* instance;
struct ll_queue* senders_list;
struct ll_queue* nodes;
struct NODEINFO_Q* local_node;
struct route_ping_pending* ping_pending;
uint64_t next_ping_req_id;
};
/**

78
src/route_node.c

@ -25,6 +25,23 @@
* (NULL если подсетей нет)
* @return количество подсетей (>= 0) или -1 при ошибке
*/
int get_node_v4_sockets(struct NODEINFO_Q *node, const struct NODEINFO_IPV4_SOCKET **out_sockets) {
if (!node || !out_sockets) {
return -1;
}
*out_sockets = NULL;
const struct NODEINFO *info = &node->node;
if (info->local_v4_sockets == 0) {
return 0;
}
const uint8_t *dynamic = (const uint8_t *)&node->node + sizeof(struct NODEINFO);
dynamic += info->node_name_len;
*out_sockets = (const struct NODEINFO_IPV4_SOCKET *)dynamic;
DEBUG_TRACE(DEBUG_CATEGORY_ROUTING, "get_node_v4_sockets: returned %u IPv4 sockets from node %p",
(unsigned)info->local_v4_sockets, (void*)node);
return (int)info->local_v4_sockets;
}
int get_node_routes(struct NODEINFO_Q *node, const struct NODEINFO_IPV4_SUBNET **out_subnets) {
if (!node || !out_subnets) {
return -1;
@ -74,7 +91,16 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG
if (s->ip.family == AF_INET) vc++;
s = s->next;
}
size_t dyn = name_len + vc * sizeof(struct NODEINFO_IPV4_SUBNET);
int sock_count = 0;
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock) {
if (e_sock->local_addr.ss_family == AF_INET) sock_count++;
e_sock = e_sock->next;
}
size_t dyn = name_len + sock_count * sizeof(struct NODEINFO_IPV4_SOCKET) + vc * sizeof(struct NODEINFO_IPV4_SUBNET);
if (!bgp->local_node) {
bgp->local_node = u_calloc(1, sizeof(struct NODEINFO_Q) + dyn);
if (!bgp->local_node) return -1;
@ -83,12 +109,20 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG
bgp->local_node->node.ver = 1;
bgp->local_node->dirty = 1;
bgp->local_node->last_ver = 1;
bgp->local_node->node.local_v4_sockets = sock_count;
bgp->local_node->node.local_v4_subnets = vc;
bgp->local_node->node.node_name_len = name_len;
memcpy(bgp->local_node->node.public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE);
}
int changed = (vc != (int)bgp->local_node->node.local_v4_subnets) || (name_len != bgp->local_node->node.node_name_len);
int changed = (vc != (int)bgp->local_node->node.local_v4_subnets)
|| (sock_count != (int)bgp->local_node->node.local_v4_sockets)
|| (name_len != bgp->local_node->node.node_name_len)
|| (memcmp(bgp->local_node->node.public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE) != 0);
if (!changed && vc > 0) {
uint8_t* current = (uint8_t*)&bgp->local_node->node + sizeof(struct NODEINFO);
current += name_len + sock_count * sizeof(struct NODEINFO_IPV4_SOCKET);
struct NODEINFO_IPV4_SUBNET* ra = (struct NODEINFO_IPV4_SUBNET*)current;
s = instance->config->my_subnets;
bool same = true;
@ -104,11 +138,39 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG
}
changed = !same;
}
if (!changed && sock_count > 0) {
uint8_t* current = (uint8_t*)&bgp->local_node->node + sizeof(struct NODEINFO);
current += name_len;
struct NODEINFO_IPV4_SOCKET* sa = (struct NODEINFO_IPV4_SOCKET*)current;
e_sock = instance->etcp_sockets;
bool same = true;
while (e_sock) {
if (e_sock->local_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&e_sock->local_addr;
if (memcmp(sa->addr, &sin->sin_addr.s_addr, 4) != 0 || sa->port != ntohs(sin->sin_port)) {
same = false;
break;
}
sa++;
}
e_sock = e_sock->next;
}
changed = !same;
}
if (changed) {
if (bgp->local_node) u_free(bgp->local_node);
bgp->local_node = u_calloc(1, sizeof(struct NODEINFO_Q) + dyn);
if (!bgp->local_node) return -1;
bgp->local_node->node.node_id = instance->node_id;
bgp->local_node->node.hop_count = 0;
uint8_t oldv = bgp->local_node->node.ver;
bgp->local_node->node.ver = ((oldv + 1) % 255) + 1;
bgp->local_node->node.local_v4_sockets = sock_count;
bgp->local_node->node.local_v4_subnets = vc;
bgp->local_node->node.node_name_len = name_len;
memcpy(bgp->local_node->node.public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE);
bgp->local_node->dirty = 1;
bgp->local_node->last_ver = bgp->local_node->node.ver;
uint8_t* dp = (uint8_t*)&bgp->local_node->node + sizeof(struct NODEINFO);
@ -116,6 +178,18 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG
memcpy(dp, instance->name, name_len);
dp += name_len;
}
struct NODEINFO_IPV4_SOCKET* sa = (struct NODEINFO_IPV4_SOCKET*)dp;
e_sock = instance->etcp_sockets;
while (e_sock) {
if (e_sock->local_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&e_sock->local_addr;
memcpy(sa->addr, &sin->sin_addr.s_addr, 4);
sa->port = ntohs(sin->sin_port);
sa++;
}
e_sock = e_sock->next;
}
dp = (uint8_t*)sa;
struct NODEINFO_IPV4_SUBNET* ra = (struct NODEINFO_IPV4_SUBNET*)dp;
s = instance->config->my_subnets;
while (s) {

13
src/route_node.h

@ -7,6 +7,7 @@
#include "secure_channel.h"
struct ROUTE_BGP;
struct ETCP_SOCKET;
/**
* @brief Информация о узле
@ -62,6 +63,9 @@ struct NODEINFO_Q {
struct ll_queue* paths; // сюда помещаем struct NODEINFO_PATH
uint8_t dirty;
uint8_t last_ver;
uint64_t last_ping_time; // время последнего замера в 0.1ms
uint16_t last_rtt; // лучший RTT в 0.1ms
struct ETCP_SOCKET* best_socket;
struct NODEINFO node; // Всегда в конце структуры - динамически расширяемый блок
};// __attribute__((packed));
@ -90,4 +94,13 @@ int route_bgp_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BG
*/
int get_node_routes(struct NODEINFO_Q *node, const struct NODEINFO_IPV4_SUBNET **out_subnets);
/**
* @brief Получает указатель на массив IPv4-сокетов узла.
*
* @param node Указатель на NODEINFO_Q
* @param out_sockets [out] указатель на первый элемент массива NODEINFO_IPV4_SOCKET
* @return количество сокетов (>= 0) или -1 при ошибке
*/
int get_node_v4_sockets(struct NODEINFO_Q *node, const struct NODEINFO_IPV4_SOCKET **out_sockets);
#endif // ROUTE_NODE_H

421
src/route_ping.c

@ -0,0 +1,421 @@
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <arpa/inet.h>
#include "../lib/platform_compat.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "utun_instance.h"
#include "etcp_api.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "route_node.h"
#include "route_bgp.h"
#include "route_ping.h"
struct route_ping_sock_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;
uint64_t request_id;
struct NODEINFO_Q* target_node;
uint8_t socket_count;
uint8_t completed_count;
uint16_t timeout_ms;
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;
uint8_t pubkey[SC_PUBKEY_SIZE];
uint8_t count_total;
uint8_t count_sent;
uint8_t count_ok;
uint32_t rtt_sum;
uint16_t timeout_ms;
uint16_t interval_ms;
void* next_timer;
};
struct route_ping_pending {
struct route_ping_pending* next;
struct ROUTE_BGP* bgp;
uint64_t request_id;
route_ping_callback_t callback;
void* arg;
void* timeout_timer;
};
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) {
struct route_ping_pending* p = (struct route_ping_pending*)arg;
if (!p) return;
p->timeout_timer = NULL;
struct ROUTE_BGP* bgp = p->bgp;
struct route_ping_pending** cur = &bgp->ping_pending;
while (*cur) {
if (*cur == p) {
*cur = p->next;
break;
}
cur = &(*cur)->next;
}
if (p->callback) {
p->callback(0, 0, 0, 0, p->arg);
}
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) {
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, "route_ping_finish: 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;
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, "route_ping_send_resp: 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,
uint8_t count, uint16_t interval_ms, uint16_t timeout_ms,
uint16_t wait_timeout_ms, route_ping_callback_t cb, void* arg) {
if (!bgp || !to_conn || count == 0 || timeout_ms == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_send_req: invalid args");
return -1;
}
struct BGP_PING_REQUEST* req_pkt = u_calloc(1, sizeof(struct BGP_PING_REQUEST));
if (!req_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_send_req: alloc failed");
return -2;
}
req_pkt->cmd = ETCP_ID_ROUTE_ENTRY;
req_pkt->subcmd = ROUTE_SUBCMD_PING_REQ;
req_pkt->request_id = bgp->next_ping_req_id++;
req_pkt->node_id = node_id;
req_pkt->count = count;
req_pkt->interval_ms = interval_ms;
req_pkt->timeout_ms = timeout_ms;
struct ll_entry* e = queue_entry_new(0);
if (!e) {
u_free(req_pkt);
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_send_req: queue_entry_new failed");
return -3;
}
e->dgram = (uint8_t*)req_pkt;
e->len = sizeof(struct BGP_PING_REQUEST);
int ret = etcp_send(to_conn, e);
if (ret != 0) {
u_free(req_pkt);
queue_entry_free(e);
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_send_req: etcp_send failed");
return -4;
}
struct route_ping_pending* pending = u_calloc(1, sizeof(struct route_ping_pending));
if (!pending) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_send_req: pending alloc failed");
return -5;
}
pending->bgp = bgp;
pending->request_id = req_pkt->request_id;
pending->callback = cb;
pending->arg = arg;
pending->next = bgp->ping_pending;
bgp->ping_pending = pending;
pending->timeout_timer = uasync_set_timeout(bgp->instance->ua, wait_timeout_ms * 10, pending, route_ping_pending_timeout);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "route_ping_send_req: request_id=%016llx node=%016llx count=%u interval=%u timeout=%u wait=%u",
(unsigned long long)pending->request_id, (unsigned long long)node_id,
(unsigned)count, (unsigned)interval_ms, (unsigned)timeout_ms, (unsigned)wait_timeout_ms);
return 0;
}
void route_ping_destroy_pending(struct ROUTE_BGP* bgp) {
if (!bgp) return;
while (bgp->ping_pending) {
struct route_ping_pending* p = bgp->ping_pending;
bgp->ping_pending = p->next;
if (p->timeout_timer) {
uasync_cancel_timeout(bgp->instance->ua, p->timeout_timer);
p->timeout_timer = NULL;
}
u_free(p);
}
}
void route_ping_handle_req(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_REQUEST)) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_handle_req: invalid args");
return;
}
const struct BGP_PING_REQUEST* req_pkt = (const struct BGP_PING_REQUEST*)data;
struct NODEINFO_Q* target = route_bgp_get_node(bgp, req_pkt->node_id);
if (!target) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "route_ping_handle_req: 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;
}
const struct NODEINFO_IPV4_SOCKET* sockets;
int sock_count = get_node_v4_sockets(target, &sockets);
if (sock_count <= 0 || !sockets) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "route_ping_handle_req: no IPv4 sockets for node %016llx",
(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;
}
// считаем подходящие локальные сокеты AF_INET
uint8_t local_count = 0;
struct ETCP_SOCKET* ls = bgp->instance->etcp_sockets;
while (ls) {
if (ls->local_addr.ss_family == AF_INET) local_count++;
ls = ls->next;
}
if (local_count == 0) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "route_ping_handle_req: no local IPv4 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, "route_ping_handle_req: alloc failed");
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;
// заполняем target_addr из первого сокета узла (можно расширить на все)
struct sockaddr_storage target_addr;
memset(&target_addr, 0, sizeof(target_addr));
struct sockaddr_in* sin = (struct sockaddr_in*)&target_addr;
sin->sin_family = AF_INET;
memcpy(&sin->sin_addr, sockets[0].addr, 4);
sin->sin_port = htons(sockets[0].port);
ls = bgp->instance->etcp_sockets;
uint8_t idx = 0;
while (ls) {
if (ls->local_addr.ss_family == AF_INET) {
req->sockets[idx].local_sock = ls;
struct route_ping_series* series = u_calloc(1, sizeof(struct route_ping_series));
if (series) {
series->req = req;
series->sock_ctx = &req->sockets[idx];
series->local_sock = ls;
series->target_addr = target_addr;
memcpy(series->pubkey, target->node.public_key, 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, target->node.public_key,
&target_addr, req_pkt->timeout_ms, route_ping_cb, series, NULL, 0);
if (ret != 0) {
// сразу считаем fail
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;
}
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 (!bgp || !from_conn || !data || len < sizeof(struct BGP_PING_RESPONSE)) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_ping_handle_resp: invalid args");
return;
}
const struct BGP_PING_RESPONSE* resp = (const struct BGP_PING_RESPONSE*)data;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "route_ping_handle_resp: 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;
p->callback(success, resp->avg_rtt, resp->count_sent, resp->count_ok, p->arg);
}
u_free(p);
return;
}
cur = &(*cur)->next;
}
DEBUG_WARN(DEBUG_CATEGORY_BGP, "route_ping_handle_resp: request_id=%016llx not found in pending list",
(unsigned long long)resp->request_id);
}

45
src/route_ping.h

@ -0,0 +1,45 @@
#ifndef ROUTE_PING_H
#define ROUTE_PING_H
#include <stdint.h>
#include <stddef.h>
#include "route_bgp.h"
#define ROUTE_SUBCMD_PING_REQ 0x07
#define ROUTE_SUBCMD_PING_RESP 0x08
struct BGP_PING_REQUEST {
uint8_t cmd; // ETCP_ID_ROUTE_ENTRY
uint8_t subcmd; // ROUTE_SUBCMD_PING_REQ
uint64_t request_id; // для корреляции
uint64_t node_id; // целевой узел (должен быть известен в nodes)
uint8_t count; // число пингов
uint16_t interval_ms; // интервал между пингами
uint16_t timeout_ms; // таймаут одного пинга
} __attribute__((packed));
struct BGP_PING_RESPONSE {
uint8_t cmd;
uint8_t subcmd; // ROUTE_SUBCMD_PING_RESP
uint64_t request_id;
uint8_t count_sent;
uint8_t count_ok;
uint16_t avg_rtt; // средний RTT в 0.1ms
uint8_t reserved[4];
} __attribute__((packed));
typedef void (*route_ping_callback_t)(int success, uint16_t avg_rtt, uint8_t count_sent, uint8_t count_ok, void* arg);
// Отправить запрос пинга через BGP, ожидать ответа с таймаутом
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,
uint16_t wait_timeout_ms, route_ping_callback_t cb, void* arg);
// Очистить список ожидающих запросов (при уничтожении BGP)
void route_ping_destroy_pending(struct ROUTE_BGP* bgp);
// Обработчики, вызываемые из route_bgp_receive_cbk
void route_ping_handle_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len);
void route_ping_handle_resp(struct ROUTE_BGP* bgp, struct ETCP_CONN* from_conn, const uint8_t* data, size_t len);
#endif // ROUTE_PING_H

6
tests/Makefile.am

@ -23,6 +23,7 @@ check_PROGRAMS = \
test_bgp_route_exchange \
test_routing_mesh \
test_etcp_ping \
test_route_ping \
bench_timeout_heap \
bench_uasync_timeouts
@ -81,6 +82,7 @@ ETCP_FULL_OBJS = \
$(top_builddir)/src/utun-config_updater.o \
$(top_builddir)/src/utun-route_lib.o \
$(top_builddir)/src/utun-route_bgp.o \
$(top_builddir)/src/utun-route_ping.o \
$(top_builddir)/src/utun-route_node.o \
$(top_builddir)/src/utun-routing.o \
$(top_builddir)/src/utun-tun_if.o \
@ -173,6 +175,10 @@ test_etcp_ping_SOURCES = test_etcp_ping.c
test_etcp_ping_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source
test_etcp_ping_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_route_ping_SOURCES = test_route_ping.c
test_route_ping_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tinycrypt/lib/include -I$(top_srcdir)/tinycrypt/lib/source
test_route_ping_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_ll_queue_SOURCES = test_ll_queue.c
test_ll_queue_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_ll_queue_LDADD = $(COMMON_LIBS)

10
tests/test_etcp_ping.c

@ -128,11 +128,15 @@ static void test_timeout(void* arg) {
}
}
static void ping_callback(int success, void* arg, uint64_t nonce) {
static void ping_callback(int success, uint16_t rtt, void* arg, uint64_t nonce,
const uint8_t* resp_data, size_t resp_data_len) {
(void)rtt;
(void)arg;
(void)nonce;
(void)resp_data;
(void)resp_data_len;
pong_received = success;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "ping_callback: success=%d", success);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "ping_callback: success=%d rtt=%u", success, (unsigned)rtt);
test_completed = 1;
}
@ -180,7 +184,7 @@ int main(void) {
struct ETCP_LINK* link = conn->links;
uint8_t peer_pubkey[SC_PUBKEY_SIZE];
memcpy(peer_pubkey, link->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE);
ret = etcp_send_ping(client_instance, peer_pubkey, &link->remote_addr, 500, ping_callback, NULL);
ret = etcp_send_ping(client_instance, peer_pubkey, &link->remote_addr, 500, ping_callback, NULL, NULL, 0);
if (ret == 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "ping sent via real link");
test_completed = 0;

Loading…
Cancel
Save