Browse Source

Заменил find_link_index на ll_queue с хеш-индексом: универсальный ключ port(2)+addr(16)+family(1)=19 байт, IPv4 дополнен нулями

topo_upd
Evgeny 2 months ago
parent
commit
ccd35074d8
  1. 162
      src/etcp_connections.c
  2. 15
      src/etcp_connections.h
  3. 7
      src/utun_instance.c

162
src/etcp_connections.c

@ -14,7 +14,6 @@
#include <string.h> #include <string.h>
#include "utun_instance.h" #include "utun_instance.h"
#include "config_parser.h" #include "config_parser.h"
#include "crc32.h"
#include "etcp.h" #include "etcp.h"
#include "stcp_link.h" #include "stcp_link.h"
#include "topo_node.h" #include "topo_node.h"
@ -345,94 +344,52 @@ restart_timer:
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive"); link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->ka_period_ms * 10, link, keepalive_timer_cb, "link_keepalive");
} }
static uint32_t sockaddr_hash(struct sockaddr_storage* addr) { static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR_KEY_SIZE]) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); memset(key, 0, LINK_ADDR_KEY_SIZE);
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6); if (!addr) return;
return crc32_calc((void*)addr, addr_len); if (addr->ss_family == AF_INET) {
} struct sockaddr_in* sa = (struct sockaddr_in*)addr;
memcpy(key, &sa->sin_port, 2);
// Бинарный поиск линка по ip_port_hash memcpy(key + 2, &sa->sin_addr.s_addr, 4);
static int find_link_index(struct ETCP_SOCKET* e_sock, uint32_t hash) { key[18] = 4;
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); } else {
if (!e_sock || e_sock->num_channels == 0) return -1; struct sockaddr_in6* sa = (struct sockaddr_in6*)addr;
memcpy(key, &sa->sin6_port, 2);
int left = 0; memcpy(key + 2, &sa->sin6_addr, 16);
int right = e_sock->num_channels - 1; key[18] = 6;
while (left <= right) {
int mid = left + (right - left) / 2;
if (e_sock->links[mid]->ip_port_hash == hash) {
return mid;
} else if (e_sock->links[mid]->ip_port_hash < hash) {
left = mid + 1;
} else {
right = mid - 1;
}
} }
return -(left + 1);
} }
// Реалокация массива линков с увеличением в 2 раза // find_link_index, realloc_links, insert_link, remove_link — УДАЛЕНЫ. Заменены на ll_queue с хеш-индексом.
static int realloc_links(struct ETCP_SOCKET* e_sock) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
size_t new_max = e_sock->max_channels == 0 ? 8 : e_sock->max_channels * 2;
struct ETCP_LINK** new_links = u_realloc(e_sock->links, new_max * sizeof(struct ETCP_LINK*));
if (!new_links) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "realloc failed");
return -1;
}
e_sock->links = new_links;
e_sock->max_channels = new_max;
return 0;
}
// Вставка линка в отсортированный массив static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) {
static int insert_link(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!e_sock || !link) return -1; if (!e_sock || !link || !e_sock->links_queue) return -1;
if (e_sock->num_channels >= e_sock->max_channels) { uint8_t key[LINK_ADDR_KEY_SIZE];
if (realloc_links(e_sock) < 0) return -1; sockaddr_to_key(&link->remote_addr, key);
} if (queue_find_data_by_index(e_sock->links_queue, key)) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "insert_link_queue: DUP addr in [%s]", e_sock->name);
int idx = find_link_index(e_sock, link->ip_port_hash);
if (idx >= 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "insert_link DUP: hash=0x%08x idx=%d existing_conn=%s existing_etcp=%s num_channels=%zu sock=%s",
link->ip_port_hash, idx,
e_sock->links[idx]->etcp ? e_sock->links[idx]->etcp->log_name : "?",
e_sock->links[idx]->etcp ? "[...]" : "?",
e_sock->num_channels, e_sock->name);
return -1; return -1;
} }
idx = -(idx + 1); struct ll_entry* qe = queue_entry_new(sizeof(struct link_queue_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "queue_entry_new failed"); return -1; }
if (idx < (int)e_sock->num_channels) { struct link_queue_entry* lqe = (struct link_queue_entry*)qe->data;
memmove(&e_sock->links[idx + 1], &e_sock->links[idx], memcpy(lqe->key, key, LINK_ADDR_KEY_SIZE);
(e_sock->num_channels - idx) * sizeof(struct ETCP_LINK*)); lqe->link = link;
} link->link_queue_entry = qe;
queue_data_put_with_index(e_sock->links_queue, qe);
e_sock->links[idx] = link;
e_sock->num_channels++;
return 0; return 0;
} }
// Удаление линка из массива static void remove_link_from_queue(struct ETCP_LINK* link) {
static void remove_link(struct ETCP_SOCKET* e_sock, uint32_t hash) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!e_sock || e_sock->num_channels == 0) return; if (!link || !link->link_queue_entry) return;
if (!link->conn || !link->conn->links_queue) return;
int idx = find_link_index(e_sock, hash); queue_remove_data(link->conn->links_queue, link->link_queue_entry);
if (idx < 0) return; queue_entry_free(link->link_queue_entry);
link->link_queue_entry = NULL;
if (idx < (int)e_sock->num_channels - 1) {
memmove(&e_sock->links[idx], &e_sock->links[idx + 1],
(e_sock->num_channels - idx - 1) * sizeof(struct ETCP_LINK*));
}
e_sock->num_channels--;
} }
static int sockaddr_equal(const struct sockaddr_storage* a, const struct sockaddr_storage* b) { static int sockaddr_equal(const struct sockaddr_storage* a, const struct sockaddr_storage* b) {
@ -450,15 +407,15 @@ static int sockaddr_equal(const struct sockaddr_storage* a, const struct sockadd
return 0; return 0;
} }
// надо править, используй sockaddr_hash
struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr) { struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!e_sock || !addr) return NULL; if (!e_sock || !addr || !e_sock->links_queue) return NULL;
uint8_t key[LINK_ADDR_KEY_SIZE];
int idx = find_link_index(e_sock, sockaddr_hash(addr)); sockaddr_to_key(addr, key);
if (idx < 0) return NULL; struct ll_entry* e = queue_find_data_by_index(e_sock->links_queue, key);
if (!e) return NULL;
return e_sock->links[idx]; struct link_queue_entry* lqe = (struct link_queue_entry*)e->data;
return lqe->link;
} }
struct ETCP_LINK* etcp_link_find_by_remote_id(struct ETCP_CONN* conn, uint8_t remote_link_id) { struct ETCP_LINK* etcp_link_find_by_remote_id(struct ETCP_CONN* conn, uint8_t remote_link_id) {
@ -710,6 +667,13 @@ struct ETCP_SOCKET* etcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_S
e_sock->nat_type = NAT_TYPE_UNKNOWN; // только результат детекции NAT, серверные сокеты стартуют с unknown e_sock->nat_type = NAT_TYPE_UNKNOWN; // только результат детекции NAT, серверные сокеты стартуют с unknown
e_sock->mtu = mtu; e_sock->mtu = mtu;
e_sock->only_local = only_local; e_sock->only_local = only_local;
e_sock->links_queue = queue_new(instance->ua, 256, offsetof(struct link_queue_entry, key), LINK_ADDR_KEY_SIZE, "links");
if (!e_sock->links_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to create links_queue for socket %s", e_sock->name);
socket_close_wrapper(e_sock->fd);
u_free(e_sock);
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "Add Socket type=%d", type); DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "Add Socket type=%d", type);
e_sock->next = instance->etcp_sockets; e_sock->next = instance->etcp_sockets;
instance->etcp_sockets = e_sock; instance->etcp_sockets = e_sock;
@ -745,12 +709,13 @@ void etcp_socket_remove(struct ETCP_SOCKET* conn) {
} }
size_t i = 0; struct ll_entry* entry;
while (i < conn->num_channels) { while ((entry = conn->links_queue->head) != NULL) {
struct ETCP_LINK* l = conn->links[i]; struct link_queue_entry* lqe = (struct link_queue_entry*)entry->data;
etcp_link_close(l); // remove_link inside shifts elements left → next at same i etcp_link_close(lqe->link);
} }
u_free(conn->links); queue_free(conn->links_queue);
conn->links_queue = NULL;
u_free(conn); u_free(conn);
} }
@ -829,13 +794,11 @@ struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn
memcpy(&link->remote_addr, remote_addr, sizeof(struct sockaddr_storage)); memcpy(&link->remote_addr, remote_addr, sizeof(struct sockaddr_storage));
link->ip_port_hash = sockaddr_hash(remote_addr);
link->last_recv_local_time = get_time_tb(); // Initialize to prevent immediate timeout link->last_recv_local_time = get_time_tb(); // Initialize to prevent immediate timeout
link->total_retransmissions = 0; link->total_retransmissions = 0;
// insert_link(conn, link); if (insert_link_queue(conn, link) < 0) {
if (insert_link(conn, link) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Can not insert link to socket"); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Can not insert link to socket");
u_free(link); u_free(link);
return NULL; return NULL;
@ -915,7 +878,7 @@ void etcp_link_close(struct ETCP_LINK* link) {
} }
if (link->etcp->last_rr_link == link) link->etcp->last_rr_link = NULL; if (link->etcp->last_rr_link == link) link->etcp->last_rr_link = NULL;
remove_link(link->conn, link->ip_port_hash); remove_link_from_queue(link);
etcp_conn_on_inflight_lim_changed(link->etcp); etcp_conn_on_inflight_lim_changed(link->etcp);
u_free(link->bbr); u_free(link->bbr);
@ -1747,11 +1710,10 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
link = existing_link; link = existing_link;
if (!sockaddr_equal(&link->remote_addr, &addr)) { if (!sockaddr_equal(&link->remote_addr, &addr)) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] IP:port changed for remote_link_id=%d socket:[%s]", conn->log_name, req->link_id, e_sock->name); DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] IP:port changed for remote_link_id=%d socket:[%s]", conn->log_name, req->link_id, e_sock->name);
if (link->conn) remove_link(link->conn, link->ip_port_hash); // remove old connection from old socket remove_link_from_queue(link); // remove from old socket queue
link->conn=e_sock; link->conn = e_sock;
memcpy(&link->remote_addr, &addr, sizeof(addr)); memcpy(&link->remote_addr, &addr, sizeof(addr));
link->ip_port_hash = sockaddr_hash(&addr); if (insert_link_queue(link->conn, link) < 0) {
if (insert_link(link->conn, link) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to reinsert link after addr change"); DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to reinsert link after addr change");
goto ec_fr; goto ec_fr;
} }
@ -2038,8 +2000,8 @@ int init_sockets(struct UTUN_INSTANCE* instance) {
} }
success_count++; success_count++;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized server %s on %s (links: %zu)", DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized server %s on %s (links: %d)",
server->name, sockaddr_storage_to_str(&server->ip).str, e_sock->num_channels); server->name, sockaddr_storage_to_str(&server->ip).str, queue_entry_count(e_sock->links_queue));
server = server->next; server = server->next;
} }

15
src/etcp_connections.h

@ -12,6 +12,7 @@ extern "C" {
#include "utun_instance.h" #include "utun_instance.h"
#include "etcp_bbr.h" #include "etcp_bbr.h"
#include "../lib/socket_compat.h" #include "../lib/socket_compat.h"
#include "../lib/ll_queue.h"
#include <stdint.h> #include <stdint.h>
#include <stddef.h> #include <stddef.h>
@ -113,6 +114,13 @@ struct PING_CONTEXT {
size_t user_data_len; size_t user_data_len;
}; };
#define LINK_ADDR_KEY_SIZE 19 // port(2) + addr(16) + family(1)
struct link_queue_entry {
uint8_t key[LINK_ADDR_KEY_SIZE];
struct ETCP_LINK* link;
};
// список активных подключений которые обслуживает сокет. каждый сокет может обслуживать много подключений // список активных подключений которые обслуживает сокет. каждый сокет может обслуживать много подключений
struct ETCP_SOCKET { struct ETCP_SOCKET {
struct ETCP_SOCKET* next; // Linked list для всех соединений struct ETCP_SOCKET* next; // Linked list для всех соединений
@ -122,10 +130,6 @@ struct ETCP_SOCKET {
struct sockaddr_storage local_addr; // Локальный адрес struct sockaddr_storage local_addr; // Локальный адрес
int mtu; // MTU для этого сокета int mtu; // MTU для этого сокета
// для входящих подключений (links) - массив упорядоченный по ip_port_hash
size_t max_channels; // сколько выделено памяти
size_t num_channels; // сколько активно
struct ETCP_LINK** links;// массив указателей на линки, сортированный по ip_port_hash
int errorcode; int errorcode;
size_t pkt_format_errors; size_t pkt_format_errors;
@ -138,6 +142,7 @@ struct ETCP_SOCKET {
uint8_t only_local; // 1 = only local connections, no forwarding uint8_t only_local; // 1 = only local connections, no forwarding
struct sockaddr_storage interface_addr; // адрес интерфейса: если bind на интерфейс - его IP, если 0.0.0.0 без интерфейса - IP интерфейса через который идет default route, иначе - адрес из конфига struct sockaddr_storage interface_addr; // адрес интерфейса: если bind на интерфейс - его IP, если 0.0.0.0 без интерфейса - IP интерфейса через который идет default route, иначе - адрес из конфига
struct sockaddr_storage nat_addr; // NAT адрес (network byte order), ss_family=0 = не определён struct sockaddr_storage nat_addr; // NAT адрес (network byte order), ss_family=0 = не определён
struct ll_queue* links_queue; // хеш-очередь линков, ключ = link_queue_entry.key (19 байт)
}; };
// NAT check status // NAT check status
@ -171,7 +176,7 @@ typedef ssize_t (*etcp_udp_send_fn_t)(socket_t fd, const void* buf, size_t len,
// ETCP Link - одно динамическое соединение (один путь) // ETCP Link - одно динамическое соединение (один путь)
struct ETCP_LINK { struct ETCP_LINK {
uint32_t ip_port_hash; // crc32 для быстрого поиска struct ll_entry* link_queue_entry; // элемент в socket->links_queue
struct ETCP_LINK* next; // Linked list подключений для ETCP_CONN (каждое подключение это child для ETCP_CONN) struct ETCP_LINK* next; // Linked list подключений для ETCP_CONN (каждое подключение это child для ETCP_CONN)
struct ETCP_CONN* etcp; // подключение (parent) struct ETCP_CONN* etcp; // подключение (parent)

7
src/utun_instance.c

@ -684,12 +684,7 @@ void utun_instance_diagnose_leaks(struct UTUN_INSTANCE *instance, const char *ph
struct ETCP_SOCKET *sock = instance->etcp_sockets; struct ETCP_SOCKET *sock = instance->etcp_sockets;
while (sock) { while (sock) {
report.etcp_sockets_count++; report.etcp_sockets_count++;
// Подсчёт линков в каждом сокете report.etcp_links_count += queue_entry_count(sock->links_queue);
for (size_t i = 0; i < sock->num_channels; i++) {
if (sock->links[i]) {
report.etcp_links_count++;
}
}
sock = sock->next; sock = sock->next;
} }

Loading…
Cancel
Save