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.
2213 lines
99 KiB
2213 lines
99 KiB
#include "etcp_connections.h" |
|
#include "etcp_api.h" |
|
#include "../lib/socket_compat.h" |
|
#include "../lib/platform_compat.h" |
|
#include "../lib/getmyip.h" |
|
#ifndef _WIN32 |
|
#include <net/if.h> |
|
#include <endian.h> |
|
#else |
|
#include <winsock2.h> |
|
#endif |
|
#include <stdlib.h> |
|
#include <unistd.h> |
|
#include <string.h> |
|
#include "utun_instance.h" |
|
#include "config_parser.h" |
|
#include "crc32.h" |
|
#include "etcp.h" |
|
#include "stcp_link.h" |
|
#include "topo_node.h" |
|
#include "topo_group.h" |
|
#include "route_ping.h" |
|
#include "../lib/memory_pool.h" |
|
#include "../lib/u_async.h" |
|
#include "../lib/debug_config.h" |
|
#include "etcp_loadbalancer.h" |
|
#include <stdlib.h> |
|
#include <time.h> |
|
#include "../lib/mem.h" |
|
#include "etcp.h" |
|
|
|
// TCP server: on new incoming connection → create minimal ETCP_CONN |
|
static void tcp_server_on_link(struct stcp_link *link, void *arg) { |
|
struct UTUN_INSTANCE *inst = (struct UTUN_INSTANCE *)arg; |
|
struct ETCP_CONN *conn = etcp_connection_create(inst, NULL); |
|
if (!conn) return; |
|
conn->transport_link = link; |
|
snprintf(conn->log_name, sizeof(conn->log_name), "tcp-[%p]", (void*)link); |
|
struct etcp_cbk_entry* cbe = inst->new_conn_cbks; |
|
while (cbe) { struct etcp_cbk_entry* n = cbe->next; cbe->fn(conn, cbe->arg); cbe = n; } |
|
{ struct etcp_cbk_entry* rcb = conn->ready_cbks; while (rcb) { struct etcp_cbk_entry* n = rcb->next; rcb->fn(conn, rcb->arg); rcb = n; } } |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server new conn=%p total=%d pending=%d", |
|
(void*)conn, |
|
queue_entry_count(inst->connections), |
|
queue_entry_count(inst->connections)); |
|
} |
|
|
|
// Forward declaration |
|
void etcp_connections_read_callback_socket(socket_t sock, void* arg); |
|
static void etcp_link_remove_from_connections(struct ETCP_SOCKET* conn, struct ETCP_LINK* link); |
|
|
|
static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision); |
|
//static int etcp_link_send_reset(struct ETCP_LINK* link); |
|
static void etcp_link_init_timer_cbk(void* arg); |
|
static void etcp_link_send_keepalive(struct ETCP_LINK* link); |
|
static void keepalive_timer_cb(void* arg); |
|
static void link_stats_timer_cb(void* arg); |
|
static void burst_resp_timeout_cb(void* arg); |
|
|
|
// === Burst sender functions === |
|
|
|
void etcp_link_burst_start(struct ETCP_LINK* link) { |
|
if (!link || !link->etcp || !link->etcp->instance) return; |
|
link->burst_active = 1; |
|
link->burst_seq = 0; |
|
link->burst_count = BURST_PACKET_COUNT; |
|
link->burst_id++; |
|
link->burst_last_time_tb = get_time_tb(); |
|
link->burst_pkt_size = 0; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] burst start id=%u count=%u", link->etcp->log_name, link->burst_id, link->burst_count); |
|
// resume send queue чтобы burst-пакеты были немедленно отправлены |
|
if (link->etcp->link_ready_for_send_fn) link->etcp->link_ready_for_send_fn(link->etcp); |
|
} |
|
|
|
void etcp_link_burst_check(struct ETCP_LINK* link) { |
|
if (!link || !link->etcp) return; |
|
if (link->burst_active) return; |
|
if (link->link_status != 1) return; |
|
if (link->inflight_bytes < link->inflight_lim_bytes) return; |
|
uint64_t now = get_time_tb(); |
|
if (now - link->burst_last_time_tb < MIN_BURST_INTERVAL_TB) return; |
|
etcp_link_burst_start(link); |
|
} |
|
|
|
void etcp_link_burst_finish(struct ETCP_LINK* link) { |
|
if (!link || !link->etcp || !link->etcp->instance) return; |
|
link->burst_active = 0; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] burst end id=%u pkt_sz=%u wait_resp", link->etcp->log_name, link->burst_id, link->burst_pkt_size); |
|
link->burst_resp_timer = uasync_set_timeout(link->etcp->instance->ua, BURST_RESP_TIMEOUT_TB, link, burst_resp_timeout_cb, "burst_resp"); |
|
if (link->etcp->link_ready_for_send_fn) link->etcp->link_ready_for_send_fn(link->etcp); |
|
} |
|
|
|
static void burst_resp_timeout_cb(void* arg) { |
|
struct ETCP_LINK* link = (struct ETCP_LINK*)arg; |
|
if (!link) return; |
|
link->burst_resp_timer = NULL; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] burst resp timeout id=%u", link->etcp->log_name, link->burst_id); |
|
} |
|
|
|
#define INIT_TIMEOUT_INITIAL 500 |
|
#define INIT_TIMEOUT_MAX 50000 |
|
|
|
static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "link=%p, is_server=%d, reset=%d, collision=%d", link, link ? link->is_server : -1, reset, collision); |
|
if (!link || !link->etcp || !link->etcp->instance) return; |
|
|
|
struct ETCP_DGRAM* dgram = u_malloc(PACKET_DATA_SIZE); |
|
if (!dgram) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "malloc failed"); |
|
return; |
|
} |
|
|
|
dgram->link = link; |
|
dgram->noencrypt_len = SC_PUBKEY_ENC_SIZE; |
|
|
|
struct ETCP_INIT_REQUEST_PKT* req = (struct ETCP_INIT_REQUEST_PKT*)dgram->data; |
|
req->code = reset ? ETCP_INIT_REQUEST : ETCP_INIT_REQUEST_NOINIT; |
|
*(uint64_t*)req->node_id = htobe64(link->etcp->instance->node_id); |
|
*(uint32_t*)req->session_id = htobe32(link->etcp->session_id); |
|
*(uint16_t*)req->mtu = htobe16(link->mtu_local); |
|
*(uint16_t*)req->keepalive = htobe16(link->keepalive_interval); |
|
*(uint16_t*)req->recovery = htobe16(link->recovery_interval / 100); |
|
req->link_id = link->local_link_id; |
|
req->socket_id = link->conn ? link->conn->sock_id : 0; |
|
req->only_local = link->conn ? link->conn->only_local : 0; |
|
req->type = link->conn ? link->conn->type : CFG_SERVER_TYPE_UNKNOWN; |
|
if (link->conn && link->conn->interface_addr.ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)&link->conn->interface_addr; |
|
memcpy(req->src_ipv4, &sin->sin_addr.s_addr, 4); |
|
*(uint16_t*)req->src_port = sin->sin_port; |
|
} else { |
|
memset(req->src_ipv4, 0, 4); |
|
memset(req->src_port, 0, 2); |
|
} |
|
req->collision = collision; |
|
memcpy(req->ed25519_pubkey, link->etcp->instance->my_ed25519_pubkey, SC_PUBKEY_SIZE); |
|
size_t offset = ETCP_INIT_REQ_SIZE; |
|
|
|
// padding |
|
int s = rand() % (link->handshake_maxsize - link->handshake_minsize) + link->handshake_minsize; |
|
int s_max = (int)(link->mtu) - (int)SC_NONCE_SIZE - 3 - (int)SC_CRC32_SIZE - (int)SC_TAG_SIZE - (int)SC_PUBKEY_ENC_SIZE + (int)UDP_SC_HDR_SIZE; |
|
if (s > s_max) s = s_max; |
|
if (s < 0) s = 0; |
|
|
|
int to_add=s-offset-UDP_HDR_SIZE - UDP_SC_HDR_SIZE; |
|
if (to_add<0) to_add=0; |
|
int max_data = PACKET_DATA_SIZE - (int)sizeof(struct ETCP_DGRAM); |
|
if (offset + to_add + SC_PUBKEY_ENC_SIZE > max_data) { |
|
to_add = max_data - offset - SC_PUBKEY_ENC_SIZE; |
|
if (to_add<0) to_add=0; |
|
} |
|
|
|
for (int i=0; i<to_add; i++) dgram->data[offset++]=rand();// fill pad |
|
// padding end |
|
|
|
uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE]; |
|
random_bytes(salt, sizeof(salt)); |
|
memcpy(dgram->data + offset, salt, SC_PUBKEY_ENC_SALT_SIZE); |
|
offset += SC_PUBKEY_ENC_SALT_SIZE; |
|
|
|
uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE]; |
|
sc_obfuscate_pubkey(salt, link->etcp->crypto_ctx.peer_public_key, |
|
link->etcp->instance->my_keys.public_key, obfuscated_pubkey); |
|
memcpy(dgram->data + offset, obfuscated_pubkey, SC_PUBKEY_SIZE); |
|
offset += SC_PUBKEY_SIZE; |
|
|
|
|
|
dgram->data_len = offset; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Sending INIT request to link, node_id=%016llx, retry=%d", (unsigned long long)link->etcp->instance->node_id, link->init_retry_count); |
|
|
|
// Debug: print remote address before sending |
|
if (link->remote_addr.ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)&link->remote_addr; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] INIT sending to %s:%d, link=%p, rst_req=%d", ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port), link, reset); |
|
} else if (link->remote_addr.ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&link->remote_addr; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] INIT sending to %s:%d, link=%p, rst_req=%d", ip_to_str(&sin6->sin6_addr, AF_INET6).str, ntohs(sin6->sin6_port), link, reset); |
|
} |
|
|
|
etcp_encrypt_send(dgram); |
|
u_free(dgram); |
|
|
|
link->init_retry_count++; |
|
} |
|
|
|
static void etcp_link_init_timer_cbk(void* arg) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
struct ETCP_LINK* link = (struct ETCP_LINK*)arg; |
|
if (!link || !link->etcp || !link->etcp->instance) return; |
|
|
|
if ((link->init_retry_count % 10) == 0 && link->init_timeout < INIT_TIMEOUT_MAX) { |
|
link->init_timeout += link->init_timeout/4 +1; |
|
if (link->init_timeout > INIT_TIMEOUT_MAX) link->init_timeout = INIT_TIMEOUT_MAX; |
|
} |
|
link->init_timer = uasync_set_timeout(link->etcp->instance->ua, link->init_timeout, link, etcp_link_init_timer_cbk, "link_init"); |
|
|
|
if (link->link_state == 1) etcp_link_send_init(link,1,0);// init (with etcp reset) |
|
else etcp_link_send_init(link,0,0);// no etcp reset (reinit) |
|
} |
|
|
|
void etcp_link_restart_init_timer(struct ETCP_LINK* link) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (link->init_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timeout = INIT_TIMEOUT_INITIAL; |
|
link->init_timer = uasync_set_timeout(link->etcp->instance->ua, link->init_timeout, link, etcp_link_init_timer_cbk, "link_init"); |
|
} |
|
|
|
void etcp_link_enter_init(struct ETCP_LINK* link) {// |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!link) return; |
|
link->link_state = 1; // handshake |
|
if (link->is_server != 0) return; |
|
etcp_link_send_init(link,1,0);// init with reset |
|
etcp_link_restart_init_timer(link); |
|
} |
|
|
|
void etcp_link_enter_reinit(struct ETCP_LINK* link) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!link) return; |
|
link->link_state = 2; // reconnect |
|
etcp_on_link_down(link->etcp); |
|
if (link->is_server != 0) return; |
|
etcp_link_send_init(link,0,0);// init without reset |
|
|
|
if (link->keepalive_timer) {// keepalive заменяяется reinit запросами |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); |
|
link->keepalive_timer = NULL; |
|
} |
|
|
|
etcp_link_restart_init_timer(link); |
|
} |
|
|
|
|
|
// Send empty keepalive packet (only timestamp, no sections) |
|
static void etcp_link_send_keepalive(struct ETCP_LINK* link) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, ""); |
|
if (!link || !link->etcp || !link->etcp->instance) return; |
|
|
|
struct ETCP_DGRAM* dgram = u_malloc(sizeof(struct ETCP_DGRAM) + 4); |
|
if (!dgram) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "malloc failed"); |
|
return; |
|
} |
|
|
|
dgram->link = link; |
|
dgram->data[0] = ETCP_KEEPALIVE; |
|
dgram->data[1] = link->ka_period_ms & 0xFF; |
|
dgram->data[2] = link->ka_period_ms >> 8; |
|
dgram->data_len = 3; |
|
dgram->noencrypt_len = 0; |
|
dgram->timestamp = get_current_timestamp(); |
|
dgram->flag_up = link->recv_keepalive; |
|
|
|
link->keepalive_sent_count++; |
|
|
|
etcp_encrypt_send(dgram); |
|
u_free(dgram); |
|
} |
|
|
|
// Check if all links for an ETCP_CONN are down |
|
// Returns 1 if all links are down or no links exist, 0 otherwise |
|
static int etcp_all_links_down(struct ETCP_CONN* etcp) { |
|
if (!etcp || !etcp->links) return 1; |
|
|
|
struct ETCP_LINK* l = etcp->links; |
|
while (l) { |
|
if (l->link_status == 1) { |
|
return 0; // At least one link is up |
|
} |
|
l = l->next; |
|
} |
|
return 1; // All links are down |
|
} |
|
|
|
static void start_keepalive_timer(struct ETCP_LINK* link) { |
|
// Start keepalive timer |
|
if (link->init_timer) {// cancel init timer |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timer = NULL; |
|
} |
|
|
|
if (link->keepalive_timer == NULL) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive timer started on link %p (interval=%d ms)", link->etcp->log_name, link, link->keepalive_interval); |
|
link->keepalive_timer = uasync_set_timeout(link->etcp->instance->ua, link->keepalive_interval * 10, link, keepalive_timer_cb, "link_keepalive"); |
|
} |
|
} |
|
|
|
// Keepalive timer callback |
|
static void keepalive_timer_cb(void* arg) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_KEEPALIVE, ""); |
|
struct ETCP_LINK* link = (struct ETCP_LINK*)arg; |
|
if (!link || !link->etcp || !link->etcp->instance) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_KEEPALIVE, "KEEPALIVE NULL !!!!!!!!"); |
|
return; |
|
} |
|
|
|
link->keepalive_timer = NULL; |
|
|
|
// Check if all links are down and start recovery if needed (client only) |
|
if (link->is_server == 0 && etcp_all_links_down(link->etcp)) { |
|
DEBUG_WARN(DEBUG_CATEGORY_KEEPALIVE, "[%s] All links are down, starting recovery", link->etcp->log_name); |
|
etcp_link_enter_reinit(link);// keepalive timr после reinit не нужен |
|
return; |
|
} |
|
|
|
// Skip if link is not initialized |
|
if (!link->initialized) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Keepalive skipped - link not initialized", |
|
link->etcp->log_name); |
|
goto restart_timer; |
|
} |
|
|
|
// Check keepalive timeout |
|
uint64_t now = get_time_tb(); |
|
uint64_t timeout_units = (uint64_t)link->keepalive_timeout * 10; // ms -> 0.1ms units |
|
uint64_t elapsed = now - link->last_recv_local_time; |
|
|
|
if (elapsed > timeout_units) { |
|
if (link->recv_keepalive != 0) { |
|
link->recv_keepalive = 0; |
|
link->link_status = 0; |
|
etcp_on_link_down(link->etcp); |
|
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] Conn:%s Link down: link_id=%d ka=%d remote_ka=%d tmo: %d>%d", link->etcp->log_name, link->conn?link->conn->name:"???", link->local_link_id, link->recv_keepalive, link->remote_keepalive, elapsed, timeout_units); |
|
DEBUG_WARN(DEBUG_CATEGORY_KEEPALIVE, "[%s] Link %p (local_id=%d) recv status changed to DOWN - no packets for %llu ms", link->etcp->log_name, link, link->local_link_id, (unsigned long long)(elapsed/10)); |
|
} |
|
} |
|
|
|
// Adaptive keepalive period (only if adaptive enabled) |
|
if (link->pkt_sent_since_keepalive) |
|
link->ka_period_ms = (uint16_t)link->keepalive_interval; |
|
else if (link->etcp->instance && link->etcp->instance->config && |
|
link->etcp->instance->config->global.keepalive_adaptive) { |
|
uint32_t next = (uint32_t)link->ka_period_ms * 105 / 100 + 1; |
|
link->ka_period_ms = next > KA_PERIOD_MAX_MS ? KA_PERIOD_MAX_MS : (uint16_t)next; |
|
} |
|
link->pkt_sent_since_keepalive = 0; |
|
|
|
// Send keepalive (server stops if link lost, client always sends) |
|
if (!link->is_server || link->recv_keepalive) |
|
etcp_link_send_keepalive(link); |
|
|
|
restart_timer: |
|
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) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6); |
|
return crc32_calc((void*)addr, addr_len); |
|
} |
|
|
|
// Бинарный поиск линка по ip_port_hash |
|
static int find_link_index(struct ETCP_SOCKET* e_sock, uint32_t hash) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!e_sock || e_sock->num_channels == 0) return -1; |
|
|
|
int left = 0; |
|
int right = e_sock->num_channels - 1; |
|
|
|
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 раза |
|
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(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!e_sock || !link) return -1; |
|
|
|
if (e_sock->num_channels >= e_sock->max_channels) { |
|
if (realloc_links(e_sock) < 0) return -1; |
|
} |
|
|
|
int idx = find_link_index(e_sock, link->ip_port_hash); |
|
if (idx >= 0) return -1; |
|
|
|
idx = -(idx + 1); |
|
|
|
if (idx < (int)e_sock->num_channels) { |
|
memmove(&e_sock->links[idx + 1], &e_sock->links[idx], |
|
(e_sock->num_channels - idx) * sizeof(struct ETCP_LINK*)); |
|
} |
|
|
|
e_sock->links[idx] = link; |
|
e_sock->num_channels++; |
|
return 0; |
|
} |
|
|
|
// Удаление линка из массива |
|
static void remove_link(struct ETCP_SOCKET* e_sock, uint32_t hash) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!e_sock || e_sock->num_channels == 0) return; |
|
|
|
int idx = find_link_index(e_sock, hash); |
|
if (idx < 0) return; |
|
|
|
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) { |
|
if (!a || !b || a->ss_family != b->ss_family) return 0; |
|
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; |
|
return (sa->sin_addr.s_addr == sb->sin_addr.s_addr && sa->sin_port == sb->sin_port); |
|
} |
|
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; |
|
return (memcmp(&sa->sin6_addr, &sb->sin6_addr, 16) == 0 && sa->sin6_port == sb->sin6_port); |
|
} |
|
return 0; |
|
} |
|
|
|
// надо править, используй sockaddr_hash |
|
struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!e_sock || !addr) return NULL; |
|
|
|
int idx = find_link_index(e_sock, sockaddr_hash(addr)); |
|
if (idx < 0) return NULL; |
|
|
|
return e_sock->links[idx]; |
|
} |
|
|
|
struct ETCP_LINK* etcp_link_find_by_remote_id(struct ETCP_CONN* conn, uint8_t remote_link_id) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!conn || remote_link_id == 0) return NULL; |
|
struct ETCP_LINK* l = conn->links; |
|
while (l) { |
|
if (l->remote_link_id == remote_link_id) return l; |
|
l = l->next; |
|
} |
|
return NULL; |
|
} |
|
|
|
int etcp_find_free_local_link_id(struct ETCP_CONN* etcp) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!etcp) return -1; |
|
|
|
// Битовый массив для 256 id (32 байта * 8 бит = 256) |
|
uint8_t used_ids[32] = {1};// индекс 0 всегда занят |
|
|
|
// Помечаем занятые id |
|
struct ETCP_LINK* link = etcp->links; |
|
while (link) { |
|
used_ids[link->local_link_id >> 3] |= (1 << (link->local_link_id & 7)); |
|
link = link->next; |
|
} |
|
|
|
// Ищем первый свободный id |
|
for (int i = 0; i < 32; i++) { |
|
if (used_ids[i] != 0xFF) { |
|
// Есть свободные биты в этом байте |
|
for (int bit = 0; bit < 8; bit++) { |
|
if (!(used_ids[i] & (1 << bit))) { |
|
return (i << 3) + bit; |
|
} |
|
} |
|
} |
|
} |
|
|
|
// Все id заняты |
|
return -1; |
|
} |
|
|
|
|
|
// =============================== |
|
|
|
struct ETCP_SOCKET* etcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_SERVER* server) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!instance || !server) return NULL; |
|
|
|
struct sockaddr_storage* ip = &server->ip; |
|
uint32_t netif_index = server->netif_index; |
|
int so_mark = server->so_mark; |
|
int fib = server->fib; |
|
uint8_t type = server->type; |
|
int mtu = server->mtu ? server->mtu : instance->config->global.mtu; |
|
if (mtu == 0 || mtu > PACKET_DATA_MAX_MTU) mtu = PACKET_DATA_MAX_MTU; |
|
uint8_t only_local = server->only_local; |
|
char* name = server->name; |
|
|
|
struct ETCP_SOCKET* e_sock = u_calloc(1, sizeof(struct ETCP_SOCKET)); |
|
if (!e_sock) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "Failed to allocate connection"); |
|
return NULL; |
|
} |
|
e_sock->fd = SOCKET_INVALID; // Initialize to invalid socket |
|
|
|
if (name && name[0]) { |
|
strncpy(e_sock->name, name, MAX_CONN_NAME_LEN - 1); |
|
e_sock->name[MAX_CONN_NAME_LEN - 1] = '\0'; |
|
} else { |
|
e_sock->name[0] = '\0'; |
|
} |
|
|
|
if (ip && (server->ipv6_mode == CFG_IPV6_MODE_TEMPORARY || server->ipv6_mode == CFG_IPV6_MODE_PERMANENT)) { |
|
if (netif_index == 0) { |
|
netif_index = get_default_route_netif_index(AF_INET6); |
|
if (netif_index == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Cannot determine default route interface for IPv6 socket %s", name); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
} |
|
int temp = (server->ipv6_mode == CFG_IPV6_MODE_TEMPORARY) ? 1 : 0; |
|
uint8_t v6addr[16]; |
|
if (get_interface_ipv6_by_index(netif_index, temp, v6addr) != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to get %s IPv6 from netif_index=%u for socket %s", |
|
temp ? "temporary" : "permanent", netif_index, name); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)ip; |
|
memcpy(&sin6->sin6_addr, v6addr, 16); |
|
} |
|
|
|
int family = AF_INET; |
|
if (ip) { |
|
family = ip->ss_family; |
|
if (family != AF_INET && family != AF_INET6) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Unsupported address family: %d", family); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
} |
|
|
|
e_sock->fd = socket_create_udp(family); |
|
if (e_sock->fd == SOCKET_INVALID) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to create socket: %s", |
|
socket_strerror(socket_get_error())); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
|
|
// Строго не используем reuseaddr, даже в тестах! |
|
socket_set_reuseaddr(e_sock->fd, 0); |
|
|
|
// Increase socket buffers for high throughput |
|
socket_set_buffers(e_sock->fd, 4 * 1024 * 1024, 4 * 1024 * 1024); |
|
|
|
if (socket_set_nonblocking(e_sock->fd) != 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "Failed to set non-blocking mode"); |
|
} |
|
|
|
// Set socket mark if specified (Linux only) |
|
if (so_mark > 0) { |
|
socket_set_mark(e_sock->fd, so_mark); |
|
} |
|
|
|
// Set FIB for FreeBSD |
|
#ifdef __FreeBSD__ |
|
if (fib > 0) { |
|
if (setsockopt(e_sock->fd, SOL_SOCKET, SO_SETFIB, &fib, sizeof(fib)) < 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "Failed to set FIB %d: %s", fib, strerror(errno)); |
|
} |
|
} |
|
#endif |
|
|
|
// Bind to interface if specified (Linux only) |
|
#ifndef _WIN32 |
|
if (netif_index > 0) { |
|
char ifname[IF_NAMESIZE]; |
|
if (if_indextoname(netif_index, ifname)) { |
|
socket_bind_to_device(e_sock->fd, ifname); |
|
} |
|
} |
|
#endif |
|
|
|
// Store the local address and bind socket if provided |
|
if (ip) { |
|
memcpy(&e_sock->local_addr, ip, sizeof(struct sockaddr_storage)); |
|
|
|
// CRITICAL: Actually bind the socket to the address |
|
socklen_t addr_len = (ip->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6); |
|
if (bind(e_sock->fd, (struct sockaddr*)ip, addr_len) < 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[ETCP] Failed to bind socket to address family %d: %s", |
|
ip->ss_family, socket_strerror(socket_get_error())); |
|
if (ip->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)ip; |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[ETCP] Failed to bind to %s:%d", ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port)); |
|
} else if (ip->ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)ip; |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[ETCP] Failed to bind to %s:%d", ip_to_str(&sin6->sin6_addr, AF_INET6).str, ntohs(sin6->sin6_port)); |
|
} |
|
socket_close_wrapper(e_sock->fd); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
|
|
if (ip->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)ip; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Successfully bound socket to local address, family=AF_INET %s:%d", ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port)); |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Listen socket initialized: name=%s fd=%d addr=%s:%d", e_sock->name, e_sock->fd, ip_to_str(&sin->sin_addr, AF_INET).str, ntohs(sin->sin_port)); |
|
} else if (ip->ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)ip; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Successfully bound socket to local address, family=AF_INET6 %s:%d", ip_to_str(&sin6->sin6_addr, AF_INET6).str, ntohs(sin6->sin6_port)); |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Listen socket initialized: name=%s fd=%d addr=%s:%d", e_sock->name, e_sock->fd, ip_to_str(&sin6->sin6_addr, AF_INET6).str, ntohs(sin6->sin6_port)); |
|
} |
|
} |
|
// Определяем interface_addr: IP интерфейса или из конфига |
|
memset(&e_sock->interface_addr, 0, sizeof(e_sock->interface_addr)); |
|
if (netif_index > 0 && ip && ip->ss_family == AF_INET) { |
|
uint32_t if_ip = get_interface_ip_by_index(netif_index); |
|
if (if_ip != 0) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)&e_sock->interface_addr; |
|
sin->sin_family = AF_INET; |
|
sin->sin_addr.s_addr = if_ip; |
|
} |
|
} else if (ip) { |
|
if (ip->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)ip; |
|
if (sin->sin_addr.s_addr == 0) { |
|
// 0.0.0.0 — определяем через default route |
|
struct sockaddr_storage remote; |
|
memset(&remote, 0, sizeof(remote)); |
|
struct sockaddr_in* rem_sin = (struct sockaddr_in*)&remote; |
|
rem_sin->sin_family = AF_INET; |
|
rem_sin->sin_port = htons(53); |
|
inet_pton(AF_INET, "8.8.8.8", &rem_sin->sin_addr); |
|
if (get_outgoing_local_ip(ip, &remote, &e_sock->interface_addr) != 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "Failed to detect outgoing local IP for socket %s", e_sock->name); |
|
} |
|
} |
|
} else if (ip->ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)ip; |
|
if (memcmp(&sin6->sin6_addr, &in6addr_any, sizeof(struct in6_addr)) == 0) { |
|
// :: — определяем через default route |
|
struct sockaddr_storage remote; |
|
memset(&remote, 0, sizeof(remote)); |
|
struct sockaddr_in6* rem_sin6 = (struct sockaddr_in6*)&remote; |
|
rem_sin6->sin6_family = AF_INET6; |
|
rem_sin6->sin6_port = htons(53); |
|
inet_pton(AF_INET6, "2001:4860:4860::8888", &rem_sin6->sin6_addr); |
|
if (get_outgoing_local_ip(ip, &remote, &e_sock->interface_addr) != 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "Failed to detect outgoing local IPv6 for socket %s", e_sock->name); |
|
} |
|
} |
|
} |
|
} |
|
if (e_sock->interface_addr.ss_family == 0) { |
|
// Не удалось определить автоматически — используем адрес из конфига |
|
if (ip) memcpy(&e_sock->interface_addr, ip, sizeof(struct sockaddr_storage)); |
|
} |
|
// Копируем порт из конфига в interface_addr |
|
if (ip && e_sock->interface_addr.ss_family == ip->ss_family) { |
|
if (ip->ss_family == AF_INET) { |
|
struct sockaddr_in* sin_cfg = (struct sockaddr_in*)ip; |
|
struct sockaddr_in* sin_if = (struct sockaddr_in*)&e_sock->interface_addr; |
|
sin_if->sin_port = sin_cfg->sin_port; |
|
} else if (ip->ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6_cfg = (struct sockaddr_in6*)ip; |
|
struct sockaddr_in6* sin6_if = (struct sockaddr_in6*)&e_sock->interface_addr; |
|
sin6_if->sin6_port = sin6_cfg->sin6_port; |
|
} |
|
} |
|
|
|
char config_str[64] = "none"; |
|
char netif_str[64] = "none"; |
|
if (e_sock->local_addr.ss_family != 0) { |
|
snprintf(config_str, sizeof(config_str), "%s", sockaddr_storage_to_str(&e_sock->local_addr).str); |
|
} |
|
if (e_sock->interface_addr.ss_family != 0) { |
|
snprintf(netif_str, sizeof(netif_str), "%s", sockaddr_storage_to_str(&e_sock->interface_addr).str); |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Listen socket initialized: name=%s fd=%d config=%s, netif=%s", e_sock->name, e_sock->fd, config_str, netif_str); |
|
|
|
e_sock->instance = instance; |
|
e_sock->errorcode = 0; |
|
e_sock->pkt_format_errors = 0; |
|
e_sock->type = type; |
|
e_sock->sock_id = instance->next_socket_id++; |
|
e_sock->nat_type = NAT_TYPE_UNKNOWN; // только результат детекции NAT, серверные сокеты стартуют с unknown |
|
e_sock->mtu = mtu; |
|
e_sock->only_local = only_local; |
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Add Socket type=%d", type); |
|
e_sock->next = instance->etcp_sockets; |
|
instance->etcp_sockets = e_sock; |
|
e_sock->socket_id = uasync_add_socket_t(instance->ua, e_sock->fd, etcp_connections_read_callback_socket, NULL, NULL, e_sock); |
|
if (!e_sock->socket_id) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to register socket with uasync"); |
|
socket_close_wrapper(e_sock->fd); |
|
u_free(e_sock); |
|
return NULL; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Registered ETCP socket with uasync"); |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Socket %p registered and active", e_sock); |
|
return e_sock; |
|
} |
|
|
|
void etcp_socket_remove(struct ETCP_SOCKET* conn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!conn) return; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Removing socket %p, socket_id=%p", conn, conn->socket_id); |
|
|
|
// Remove from uasync if registered |
|
if (conn->socket_id) { |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Removing socket from uasync, instance=%p, ua=%p", conn->instance, conn->instance->ua); |
|
uasync_remove_socket_t(conn->instance->ua, conn->fd); |
|
conn->socket_id = NULL; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Unregistered socket from uasync"); |
|
} |
|
|
|
if (conn->fd != SOCKET_INVALID) { |
|
socket_close_wrapper(conn->fd); |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Closed socket"); |
|
} |
|
|
|
|
|
size_t i = 0; |
|
while (i < conn->num_channels) { |
|
struct ETCP_LINK* l = conn->links[i]; |
|
etcp_link_close(l); // remove_link inside shifts elements left → next at same i |
|
} |
|
u_free(conn->links); |
|
|
|
u_free(conn); |
|
} |
|
|
|
|
|
struct ETCP_LINK* etcp_link_new(struct ETCP_CONN* etcp, struct ETCP_SOCKET* conn, struct sockaddr_storage* remote_addr, uint8_t is_server) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!remote_addr) return NULL; |
|
|
|
struct ETCP_LINK* link = u_calloc(1, sizeof(struct ETCP_LINK)); |
|
if (!link) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "calloc failed - out of memory or pool exhausted"); |
|
return NULL; |
|
} |
|
|
|
link->conn = conn; |
|
link->etcp = etcp; |
|
link->is_server = is_server; |
|
int mtu = conn->mtu; |
|
if (mtu == 0) mtu = 1500; |
|
if (mtu > PACKET_DATA_MAX_MTU) mtu = PACKET_DATA_MAX_MTU; |
|
link->mtu_local = mtu; |
|
link->mtu = mtu; |
|
etcp_update_mtu(etcp); |
|
link->initialized = 0; |
|
link->init_timer = NULL; |
|
link->init_timeout = 0; |
|
link->init_retry_count = 0; |
|
link->link_status = 0; // down initially |
|
link->send_hook = NULL; |
|
link->send_hook_ctx = NULL; |
|
link->handshake_minsize = 100; |
|
link->handshake_maxsize = mtu;// 28 = udp header size |
|
|
|
// Initialize keepalive timeout from global config |
|
if (etcp->instance && etcp->instance->config) { |
|
link->keepalive_timeout = etcp->instance->config->global.keepalive_timeout; |
|
link->keepalive_interval = etcp->instance->config->global.keepalive_interval; |
|
} else { |
|
link->keepalive_timeout = 2000; // Default 2 seconds |
|
link->keepalive_interval = 200; // Default 0.2 s |
|
} |
|
if (link->keepalive_interval < 10) link->keepalive_interval = 10; |
|
link->keepalive_sent_count = 0; |
|
link->keepalive_recv_count = 0; |
|
link->ka_period_ms = (uint16_t)link->keepalive_interval; |
|
link->inflight_lim_bytes = link->mtu * 4; // BBR init_cwnd (~4 packets) |
|
link->bandwidth = 10000; // начальная оценка 10 Mbps для шейпера |
|
link->burst_id = 0; |
|
link->burst_active = 0; |
|
link->burst_last_time_tb = 0; |
|
link->burst_target_bdp = 0; |
|
|
|
link->delivered_bytes = 0; |
|
link->acked_bytes = 0; |
|
link->acked_packets = 0; |
|
link->last_ack_time_tb = get_time_tb(); |
|
link->bbr_pacing_rate = 0; |
|
link->bbr_loss_since_ack = 0; |
|
link->bbr = u_calloc(1, sizeof(struct bbr)); |
|
if (!link->bbr) { |
|
u_free(link); |
|
return NULL; |
|
} |
|
bbr_init(link->bbr); |
|
|
|
// Выделяем свободный local_link_id |
|
int free_id = etcp_find_free_local_link_id(etcp); |
|
if (free_id <= 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "no free local_link_id available"); |
|
u_free(link); |
|
return NULL; |
|
} |
|
link->local_link_id = (uint8_t)free_id; |
|
if (link->bbr) link->bbr->link_id = (uint8_t)free_id; |
|
|
|
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->total_retransmissions = 0; |
|
|
|
// insert_link(conn, link); |
|
if (insert_link(conn, link) < 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Can not insert link to socket"); |
|
u_free(link); |
|
return NULL; |
|
} |
|
|
|
struct ETCP_LINK* l=etcp->links; |
|
while (l && l->next) l=l->next; |
|
if (l) l->next = link; else etcp->links = link; |
|
|
|
etcp_link_update_inflight_lim(link, link->inflight_lim_bytes); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "NEW link initialized on etcp=[%s] link=%p socket=%s id=%d is_server=%d mtu=%d", etcp->log_name, link, conn->name, link->local_link_id, link->is_server, link->mtu); |
|
|
|
if (is_server == 0) { |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "client link, calling etcp_link_send_init"); |
|
etcp_link_enter_init(link); |
|
} |
|
|
|
return link; |
|
} |
|
|
|
void etcp_link_update_inflight_lim(struct ETCP_LINK* link, uint32_t new_lim) { |
|
if (!link || !link->etcp) return; |
|
if (new_lim < INFLIGHT_LIM_MIN) new_lim = INFLIGHT_LIM_MIN; |
|
if (new_lim > link->etcp->max_inflight) new_lim = link->etcp->max_inflight; |
|
|
|
link->inflight_lim_bytes = new_lim; |
|
etcp_conn_on_inflight_lim_changed(link->etcp); |
|
} |
|
|
|
void etcp_link_close(struct ETCP_LINK* link) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!link) return; |
|
|
|
if (link->burst_resp_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->burst_resp_timer); |
|
link->burst_resp_timer = NULL; |
|
} |
|
|
|
if (!link->conn) { |
|
struct ETCP_LINK **pp = &link->etcp->links; |
|
while (*pp && *pp != link) pp = &(*pp)->next; |
|
if (*pp) *pp = link->next; |
|
if (link->init_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
if (link->shaper_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->shaper_timer); |
|
if (link->keepalive_timer) uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); |
|
etcp_conn_on_inflight_lim_changed(link->etcp); |
|
u_free(link); |
|
return; |
|
} |
|
|
|
// Cancel init timer if active |
|
if (link->init_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timer = NULL; |
|
} |
|
|
|
// Cancel shaper timer if active |
|
if (link->shaper_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->shaper_timer); |
|
link->shaper_timer = NULL; |
|
} |
|
|
|
// Cancel keepalive timer if active |
|
if (link->keepalive_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->keepalive_timer); |
|
link->keepalive_timer = NULL; |
|
} |
|
|
|
// универсальное удаление из односвязного списка |
|
struct ETCP_LINK **pp = &link->etcp->links; |
|
while (*pp) { |
|
if (*pp == link) { |
|
*pp = link->next; |
|
break; |
|
} |
|
pp = &(*pp)->next; |
|
} |
|
if (link->etcp->last_rr_link == link) link->etcp->last_rr_link = NULL; |
|
|
|
remove_link(link->conn, link->ip_port_hash); |
|
|
|
etcp_conn_on_inflight_lim_changed(link->etcp); |
|
u_free(link->bbr); |
|
u_free(link); |
|
} |
|
|
|
ssize_t etcp_udp_send(struct ETCP_LINK* link, socket_t fd, const void* buf, size_t len, |
|
const struct sockaddr* addr, socklen_t addr_len) { |
|
if (link && link->send_hook) |
|
return link->send_hook(fd, buf, len, addr, addr_len, link, link->send_hook_ctx); |
|
return socket_sendto(fd, buf, len, addr, addr_len); |
|
} |
|
|
|
int etcp_encrypt_send(struct ETCP_DGRAM* dgram) { |
|
if (!dgram || !dgram->link) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Null pointer"); return -1; } |
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Send rk=%d lk=%d up=%d", |
|
dgram->link->etcp->log_name, dgram->link->recv_keepalive, dgram->link->remote_keepalive, dgram->link->link_status); |
|
|
|
// Mark that packet was sent (for keepalive logic) |
|
dgram->link->pkt_sent_since_keepalive = 1; |
|
dgram->flag_up=dgram->link->recv_keepalive; |
|
// Размер зашифрованной части не должен превышать UDP payload = MTU - заголовки |
|
int errcode=0; |
|
sc_context_t* sc = &dgram->link->etcp->crypto_ctx; |
|
int len=dgram->data_len-dgram->noencrypt_len;// не забываем добавить timestamp (2 bytes) |
|
if (len<0 || len>dgram->link->mtu - UDP_HDR_SIZE) { dgram->link->send_errors++; errcode=1; goto es_err; } |
|
uint8_t enc_buf[1600]; |
|
size_t enc_buf_len=0; |
|
dgram->timestamp=get_current_timestamp(); |
|
|
|
dgram->link->total_encrypted += dgram->data_len; |
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "Before encryption", dgram->data, dgram->data_len); |
|
sc_encrypt(sc, (uint8_t*)&dgram->timestamp, 3 + len, enc_buf, &enc_buf_len); |
|
if (enc_buf_len == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "eencryption failed for node %016llx", (unsigned long long)dgram->link->etcp->instance->node_id); |
|
dgram->link->send_errors++; errcode=2; goto es_err; } |
|
if (enc_buf_len + dgram->noencrypt_len > (size_t)(dgram->link->mtu - UDP_HDR_SIZE)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "packet too long len=%d ne_len=%d", enc_buf_len, dgram->noencrypt_len); |
|
dgram->link->send_errors++; errcode=3; goto es_err; } |
|
memcpy(enc_buf+enc_buf_len, dgram->data+len, dgram->noencrypt_len); |
|
|
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "Encrypted", enc_buf, enc_buf_len + dgram->noencrypt_len); |
|
|
|
struct sockaddr_storage* addr=&dgram->link->remote_addr; |
|
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6); |
|
|
|
// Debug: print where we're sending the packet |
|
if (addr->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)addr; |
|
// inet_ntop removed - use ip_to_str |
|
} |
|
|
|
ssize_t sent = etcp_udp_send(dgram->link, dgram->link->conn->fd, enc_buf, enc_buf_len + dgram->noencrypt_len, |
|
(struct sockaddr*)addr, addr_len); |
|
if (sent < 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "sendto failed, sock_err=%d dst=%s fd=%d", |
|
socket_get_error(), sockaddr_storage_to_str(addr).str, dgram->link->conn->fd); |
|
dgram->link->send_errors++; errcode=4; goto es_err; |
|
} |
|
return (int)sent; |
|
es_err: |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "error %d", errcode); |
|
return -1; |
|
} |
|
|
|
static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, socket_t fd, sc_context_t* sc, const struct sockaddr_storage* addr) { |
|
if (!dgram || !sc || !addr) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "Null pointer in ping send"); |
|
return -1; |
|
} |
|
int len = dgram->data_len - dgram->noencrypt_len; |
|
if (len < 0 || len > PACKET_DATA_SIZE - UDP_HDR_SIZE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet data invalid len=%d", len); |
|
return -1; |
|
} |
|
uint8_t enc_buf[1600]; |
|
size_t enc_buf_len = 0; |
|
dgram->timestamp = get_current_timestamp(); |
|
dgram->flag_up = 1; |
|
sc_encrypt(sc, (uint8_t*)&dgram->timestamp, 3 + len, enc_buf, &enc_buf_len); |
|
if (enc_buf_len == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "encryption failed for ping"); |
|
return -1; |
|
} |
|
if (enc_buf_len + dgram->noencrypt_len > (size_t)(PACKET_DATA_SIZE - UDP_HDR_SIZE)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "ping packet too long enc=%zu ne=%d", enc_buf_len, dgram->noencrypt_len); |
|
return -1; |
|
} |
|
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); |
|
ssize_t sent = enc_buf_len + dgram->noencrypt_len; |
|
// uint8_t* xaddr=&((struct sockaddr_in*)addr)->sin_addr; |
|
//xaddr[0]=192; xaddr[1]=168; xaddr[2]=10; xaddr[3]=1; |
|
sent = socket_sendto(fd, enc_buf, enc_buf_len + dgram->noencrypt_len, (struct sockaddr*)addr, addr_len); |
|
if (sent < 0) { |
|
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; |
|
} |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ping sendto succeeded to %s sent=%zd bytes fd=%d", sockaddr_storage_to_str(addr).str, sent, fd); |
|
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; |
|
if (ctx->timeout_timer) { |
|
uasync_cancel_timeout(ctx->instance->ua, ctx->timeout_timer); |
|
ctx->timeout_timer = NULL; |
|
} |
|
ctx->cb(0, 0, ctx->arg, ctx->nonce, NULL, 0); |
|
if (ctx->instance->pending_pings == ctx) { |
|
ctx->instance->pending_pings = ctx->next; |
|
} else { |
|
struct PING_CONTEXT* prev = ctx->instance->pending_pings; |
|
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_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_BGP, "bad args"); |
|
return -1; |
|
} |
|
if (user_data_len > PACKET_DATA_SIZE - 24 - 2 - SC_PUBKEY_ENC_SIZE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "user_data too long"); |
|
return -2; |
|
} |
|
struct PING_CONTEXT* ctx = u_malloc(sizeof(struct PING_CONTEXT)); |
|
if (!ctx) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "malloc ctx"); |
|
return -3; |
|
} |
|
ctx->next = NULL; |
|
ctx->instance = instance; |
|
ctx->cb = cb; |
|
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_BGP, "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_BGP, "malloc dgram"); |
|
return -4; |
|
} |
|
dgram->link = NULL; |
|
dgram->noencrypt_len = SC_PUBKEY_ENC_SIZE; |
|
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); 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(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(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); |
|
if (ctx->user_data) u_free(ctx->user_data); |
|
u_free(ctx); |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "set key failed"); |
|
return -5; |
|
} |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ping sent nonce=%016llx timeout=%d ulen=%zu", |
|
(unsigned long long)ctx->nonce, timeout_ms, user_data_len); |
|
etcp_send_ping_raw(dgram, e_sock->fd, &sc, addr); |
|
u_free(dgram); |
|
if (instance->pending_pings == NULL) { |
|
instance->pending_pings = ctx; |
|
} else { |
|
struct PING_CONTEXT* last = instance->pending_pings; |
|
while (last->next) last = last->next; |
|
last->next = ctx; |
|
} |
|
ctx->timeout_timer = uasync_set_timeout(instance->ua, timeout_ms * 10, ctx, ping_timeout_cbk, "ping_timeout"); |
|
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_BGP, "bad args: inst=%p pubkey=%p addr=%p timeout=%d cb=%p udata=%p udata_len=%zu", |
|
(void*)instance, (void*)peer_pubkey_bin, (void*)addr, timeout_ms, |
|
(void*)(uintptr_t)cb, (void*)user_data, user_data_len); |
|
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_BGP, "no socket for addr_family=%d", addr->ss_family); |
|
return -2; |
|
} |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ping N1 [%s]", sockaddr_storage_to_str(addr).str); |
|
return etcp_send_ping_to_socket(instance, e_sock, peer_pubkey_bin, addr, timeout_ms, |
|
cb, user_arg, user_data, user_data_len); |
|
} |
|
|
|
// === Helpers extracted from etcp_connections_read_callback_socket === |
|
|
|
static int handle_ping(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, const struct sockaddr_storage* addr, const uint8_t* decrypted_pubkey, size_t pkt_len) { |
|
if (pkt_len < 22) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PING too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str); |
|
return 7; |
|
} |
|
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; |
|
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); 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(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(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) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG send nonce=%016llx to=%s fd=%d", |
|
(unsigned long long)nonce, sockaddr_storage_to_str(addr).str, e_sock->fd); |
|
etcp_send_ping_raw(resp, e_sock->fd, &resp_sc, addr); |
|
} |
|
u_free(resp); |
|
} |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return 0; |
|
} |
|
|
|
static int handle_pong(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, const struct sockaddr_storage* addr, size_t pkt_len) { |
|
if (pkt_len < 20) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PONG too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str); |
|
return 7; |
|
} |
|
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; |
|
} |
|
} |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG recv nonce=%016llx data_len=%u from=%s socket=%s", |
|
(unsigned long long)nonce, (unsigned)pkt->data_len, |
|
sockaddr_storage_to_str(addr).str, e_sock->name); |
|
struct PING_CONTEXT* ctx = e_sock->instance->pending_pings; |
|
struct PING_CONTEXT* prev = NULL; |
|
int found = 0; |
|
while (ctx) { |
|
if (ctx->nonce == nonce) { |
|
found = 1; |
|
if (prev) prev->next = ctx->next; |
|
else e_sock->instance->pending_pings = ctx->next; |
|
if (ctx->timeout_timer) { |
|
uasync_cancel_timeout(e_sock->instance->ua, ctx->timeout_timer); |
|
ctx->timeout_timer = NULL; |
|
} |
|
uint64_t now = get_time_tb(); |
|
uint16_t rtt = (now >= ctx->send_time) ? (uint16_t)(now - ctx->send_time) : 0; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG matched nonce=%016llx rtt=%u", |
|
(unsigned long long)nonce, (unsigned)rtt); |
|
ctx->cb(1, rtt, ctx->arg, nonce, udata, ulen); |
|
if (ctx->user_data) u_free(ctx->user_data); |
|
u_free(ctx); |
|
break; |
|
} |
|
prev = ctx; |
|
ctx = ctx->next; |
|
} |
|
if (!found) { |
|
DEBUG_WARN(DEBUG_CATEGORY_BGP, "PONG nonce=%016llx NOT FOUND in pending (timeout?)", |
|
(unsigned long long)nonce); |
|
} |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return 0; |
|
} |
|
|
|
static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, struct ETCP_LINK* link, struct ETCP_CONN* conn, const struct ETCP_INIT_REQUEST_PKT* req, const struct sockaddr_storage* addr, uint8_t send_reset, uint32_t req_src_ip, uint16_t req_src_port, size_t pkt_len) { |
|
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data; |
|
|
|
// Set response code: 0x03 (with reset) or 0x05 (without reset) |
|
// response with init (0x03) only if reinit was actually done on server side |
|
if (send_reset != 0) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response with reset"); |
|
resp->code = ETCP_INIT_RESPONSE; // 0x03 - with reset |
|
} else { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response without reset"); |
|
resp->code = ETCP_INIT_RESPONSE_NOINIT; // 0x05 - without reset |
|
} |
|
*(uint64_t*)resp->node_id = htobe64(e_sock->instance->node_id); |
|
*(uint32_t*)resp->session_id = htobe32(conn->session_id); |
|
|
|
resp->mtu[0]=link->mtu_local>>8; |
|
resp->mtu[1]=link->mtu_local; |
|
resp->link_id = link->local_link_id; |
|
resp->remote_socket_id = req->socket_id; |
|
resp->only_local = e_sock->only_local; |
|
resp->type = e_sock->type; |
|
// Add client's IP:port (so client behind NAT can know its external address) |
|
if (addr->ss_family == AF_INET) { |
|
struct sockaddr_in *sin = (struct sockaddr_in*)addr; |
|
memcpy(resp->peer_ipv4, &sin->sin_addr.s_addr, 4); |
|
uint16_t port = ntohs(sin->sin_port); |
|
resp->peer_port[0] = port >> 8; |
|
resp->peer_port[1] = port & 0xFF; |
|
link->nat_ip = sin->sin_addr.s_addr; |
|
link->nat_port = port; |
|
} else { |
|
// For IPv6, set to 0 (not supported for NAT traversal) |
|
memset(resp->peer_ipv4, 0, 4); |
|
memset(resp->peer_port, 0, 2); |
|
link->nat_ip = 0; |
|
link->nat_port = 0; |
|
} |
|
// DIRECT detection: if client reports its own address and it matches observed source → real public IP |
|
if (pkt_len >= ETCP_INIT_REQ_SIZE && addr->ss_family == AF_INET && link->nat_ip != 0) { |
|
if (req_src_ip != 0 && req_src_ip == link->nat_ip |
|
&& req_src_port == link->nat_port |
|
&& !is_local_subnet(link->nat_ip)) |
|
{ |
|
link->nat_type = NAT_TYPE_DIRECT; |
|
link->nat_check_status = NAT_CHECK_EIM; |
|
if (link->etcp->instance->topo_groups) { |
|
topo_group_send_nat_info(link->etcp, link->remote_socket_id, |
|
link->nat_ip, link->nat_port, NAT_TYPE_DIRECT); |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "DIRECT IP: %s:%u for %s", |
|
ip_to_str(&link->nat_ip, AF_INET).str, link->nat_port, |
|
link->etcp->log_name); |
|
} |
|
} |
|
/* Self-detection: if my IP is non-local, I'm on a public IP */ |
|
DEBUG_INFO(DEBUG_CATEGORY_NAT, "self-detect check: socket=%s type=%d nat_type=%d if_addr=%s local_addr=%s", |
|
e_sock->name, e_sock->type, e_sock->nat_type, |
|
sockaddr_storage_to_str(&e_sock->interface_addr).str, |
|
sockaddr_storage_to_str(&e_sock->local_addr).str); |
|
if (e_sock->nat_type != NAT_VERIFIED_DIRECT && e_sock->interface_addr.ss_family == AF_INET) { |
|
uint32_t my_ip = ((struct sockaddr_in*)&e_sock->interface_addr)->sin_addr.s_addr; |
|
if (my_ip != 0 && !is_local_subnet(my_ip)) { |
|
e_sock->nat_type = NAT_VERIFIED_DIRECT; |
|
struct sockaddr_in* nat_sin = (struct sockaddr_in*)&e_sock->nat_addr; |
|
nat_sin->sin_family = AF_INET; |
|
nat_sin->sin_addr.s_addr = my_ip; |
|
nat_sin->sin_port = ((struct sockaddr_in*)&e_sock->interface_addr)->sin_port; |
|
{ struct TOPO_GROUP* g = topo_groups_get_default(link->etcp->instance->topo_groups); |
|
if (g) { |
|
topo_group_update_my_nodeinfo(g->instance, g); |
|
if (g->local_node && g->senders_list) { |
|
g->local_node->node->ver = (g->local_node->node->ver % 255) + 1; |
|
g->local_node->dirty = 1; |
|
struct ll_entry* se = g->senders_list->head; |
|
while (se) { |
|
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; |
|
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn); |
|
se = se->next; |
|
} |
|
} |
|
} |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_NAT, "Self-detected PUBLIC: socket=%s ip=%s port=%u sock_id=%d", |
|
e_sock->name, ip_to_str(&my_ip, AF_INET).str, |
|
ntohs(nat_sin->sin_port), e_sock->sock_id); |
|
} |
|
} |
|
memcpy(resp->ed25519_pubkey, e_sock->instance->my_ed25519_pubkey, SC_PUBKEY_SIZE); |
|
pkt->noencrypt_len=0; |
|
pkt->link=link; |
|
link->recv_keepalive = 1; |
|
link->last_recv_local_time = get_time_tb(); |
|
link->last_recv_timestamp = pkt->timestamp; |
|
|
|
int xoffset=sizeof(struct ETCP_INIT_RESPONSE_PKT); |
|
// padding |
|
int s = rand() % (link->handshake_maxsize - link->handshake_minsize) + link->handshake_minsize; |
|
if (s > (int)(link->mtu)) s = (int)(link->mtu); |
|
if (s < 0) s = 0; |
|
|
|
int to_add=s - xoffset - UDP_HDR_SIZE - UDP_SC_HDR_SIZE; |
|
if (to_add<0) to_add=0; |
|
if (xoffset + to_add > PACKET_DATA_SIZE) { to_add = PACKET_DATA_SIZE - xoffset; if (to_add<0) to_add=0; } |
|
|
|
for (int i=0; i<to_add; i++) pkt->data[xoffset++]=rand();// fill pad |
|
// padding end |
|
pkt->data_len=xoffset; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "Sending INIT RESPONSE, link=%p, local_link_id=%d, remote_link_id=%d, dst=%s fd=%d", |
|
link, link->local_link_id, link->remote_link_id, |
|
sockaddr_storage_to_str(&link->remote_addr).str, link->conn->fd); |
|
etcp_encrypt_send(pkt); |
|
|
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
link->initialized = 1; |
|
link->link_state = 3; |
|
if (link->init_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timer = NULL; |
|
} |
|
if (link->etcp->initialized == 0) { |
|
etcp_conn_ready(link->etcp); |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Connection established: log_name=%s socket=%s link_id=%d status=UP", link->etcp->log_name, e_sock->name, link->local_link_id); |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (server)", link->etcp->log_name, link, link->local_link_id); |
|
start_keepalive_timer(link); |
|
loadbalancer_link_ready(link); |
|
// Restart NAT check after link is up (e.g. after address change or reinit) |
|
{ struct TOPO_GROUP* g = topo_groups_get_default(link->etcp->instance->topo_groups); |
|
if (g && link->nat_check_status < NAT_CHECK_IN_PROGRESS) topo_group_start_link_nat_check(g, link); } |
|
} |
|
|
|
static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, struct ETCP_LINK* link, uint8_t pkt_code, size_t pkt_len) { |
|
if (pkt_len < ETCP_INIT_RESP_SIZE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE too short: pkt_len=%zu", pkt_len); |
|
return 46; |
|
} |
|
// ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN |
|
// ETCP_INIT_RESPONSE_NOINIT (0x05) - no reset |
|
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data; |
|
uint64_t server_node_id = be64toh(*(uint64_t*)resp->node_id); |
|
uint32_t resp_session_id = be32toh(*(uint32_t*)resp->session_id); |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "INIT_RESPONSE session_id=%08x", resp_session_id); |
|
// Check session_id: ignore response if it doesn't match our session |
|
if (resp_session_id != link->etcp->session_id) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] INIT_RESPONSE session_id mismatch: got %08x, expected %08x, ignoring", |
|
link->etcp->log_name, resp_session_id, link->etcp->session_id); |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return 0; |
|
} |
|
link->mtu_remote = be16toh(*(uint16_t*)resp->mtu); |
|
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU; |
|
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote; |
|
etcp_update_mtu(link->etcp); |
|
link->remote_link_id = resp->link_id; |
|
link->remote_socket_id = resp->remote_socket_id; |
|
link->remote_only_local = resp->only_local; |
|
link->remote_type = resp->type; |
|
|
|
// Parse NAT IP:port from response (new format includes 4+2 bytes) |
|
if (pkt_len >= ETCP_INIT_RESP_SIZE) { |
|
uint32_t new_nat_ip; |
|
memcpy(&new_nat_ip, resp->peer_ipv4, 4); |
|
uint16_t new_nat_port = be16toh(*(uint16_t*)resp->peer_port); |
|
|
|
// Check if NAT address changed |
|
if (link->nat_ip == 0 && link->nat_port == 0) { |
|
// First time receiving NAT info |
|
link->nat_ip = new_nat_ip; |
|
link->nat_port = new_nat_port; |
|
struct in_addr addr; |
|
addr.s_addr = new_nat_ip; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] NAT address initialized: %s:%u", |
|
link->etcp->log_name, ip_to_str(&addr, AF_INET).str, new_nat_port); |
|
} else if (link->nat_ip != new_nat_ip || link->nat_port != new_nat_port) { |
|
// NAT address changed |
|
struct in_addr old_addr, new_addr; |
|
old_addr.s_addr = link->nat_ip; |
|
new_addr.s_addr = new_nat_ip; |
|
|
|
link->nat_ip = new_nat_ip; |
|
link->nat_port = new_nat_port; |
|
link->nat_changes_count++; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] NAT address changed: %s:%u -> %s:%u (change #%u)", |
|
link->etcp->log_name, ip_to_str(&old_addr.s_addr, AF_INET).str, link->nat_port, ip_to_str(&new_addr.s_addr, AF_INET).str, new_nat_port, |
|
link->nat_changes_count); |
|
} |
|
|
|
// Update socket NAT address (only if not already verified) |
|
{ |
|
DEBUG_INFO(DEBUG_CATEGORY_NAT, "init_resp nat: socket=%s nat_type=%d new_ip=%s new_port=%u", |
|
e_sock->name, e_sock->nat_type, |
|
ip_to_str(&new_nat_ip, AF_INET).str, new_nat_port); |
|
if (e_sock->nat_type < NAT_VERIFIED_UNKNOWN) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)&e_sock->nat_addr; |
|
sin->sin_family = AF_INET; |
|
sin->sin_addr.s_addr = new_nat_ip; |
|
sin->sin_port = htons(new_nat_port); |
|
} else { |
|
DEBUG_INFO(DEBUG_CATEGORY_NAT, "init_resp nat update SKIPPED (already verified): socket=%s nat_type=%d", |
|
e_sock->name, e_sock->nat_type); |
|
} |
|
} |
|
} else { |
|
// Legacy format without NAT info |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Received legacy INIT_RESPONSE without NAT info", |
|
link->etcp->log_name); |
|
} |
|
|
|
memcpy(link->remote_ed25519_pubkey, resp->ed25519_pubkey, SC_PUBKEY_SIZE); |
|
memcpy(link->etcp->peer_ed25519_pubkey, resp->ed25519_pubkey, SC_PUBKEY_SIZE); |
|
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "[%s] Received Ed25519 pubkey from peer", link->etcp->log_name); |
|
|
|
|
|
// DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Received INIT_RESPONSE from server_node_id=%llu, mtu=%d", (unsigned long long)server_node_id, link->mtu); |
|
|
|
etcp_conn_set_peer_node_id(link->etcp, server_node_id); |
|
|
|
link->initialized = 1;// получен init response (client) |
|
link->link_state = 3; // connected |
|
if (link->init_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timer = NULL; |
|
} |
|
|
|
if (pkt_code == ETCP_INIT_RESPONSE && !link->etcp->reset_done) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT from client: INIT_RESPONSE(0x03) received, reinit conn=%p", |
|
link->etcp->log_name, link->etcp); |
|
etcp_conn_reinit(link->etcp); |
|
} |
|
|
|
if (link->etcp->initialized == 0) { |
|
etcp_conn_ready(link->etcp); |
|
} |
|
|
|
loadbalancer_link_ready(link); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (client): rk=%d lk=%d up=%d ki=%d", |
|
link->etcp->log_name, link, link->local_link_id, link->recv_keepalive, link->remote_keepalive, link->link_status, link->keepalive_interval); |
|
|
|
// Start keepalive timer |
|
etcp_link_send_keepalive(link); |
|
|
|
start_keepalive_timer(link); |
|
|
|
|
|
loadbalancer_link_ready(link); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp client: Link initialized successfully! Server node_id=%016llx, mtu=%d, local_link_id=%d, remote_link_id=%d", (unsigned long long)server_node_id, link->mtu, link->local_link_id, link->remote_link_id); |
|
|
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return 0; |
|
} |
|
|
|
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); |
|
// !!!!!! DANGER: в этой функции ПРЕДЕЛЬНАЯ АККУРАТНОСТЬ. Если кажется что не туда указатель то невнимательно аланизировал !!!!! |
|
// НЕ РУИНИТЬ (uint8_t*)&pkt->timestamp - это правильно !!!! |
|
// |
|
// Ошибки функции (errorcode): |
|
// 1 - пакет слишком маленький для init (< SC_PUBKEY_SIZE) |
|
// 2 - не удалось установить peer public key при init |
|
// 3 - не удалось расшифровать init пакет |
|
// 4 - не init/p ing пакет (неверный код) |
|
// 5 - коллизия peer ID и ключей |
|
// 6 - не удалось расшифровать обычный пакет |
|
// 7 - слишком короткий пакет |
|
// 8 - ключ не в списке allowed_keys |
|
// 13 - переполнение при парсинге пакета |
|
// 46 - расшифрованный пакет слишком маленький (< 3 байта) |
|
// 55 - не удалось создать подключение |
|
// 66 - не удалось создать линк |
|
struct ETCP_SOCKET* e_sock = (struct ETCP_SOCKET*)arg; |
|
if (!e_sock) return; |
|
|
|
struct sockaddr_storage addr; |
|
uint8_t data[PACKET_DATA_SIZE]; |
|
socklen_t addr_len=sizeof(addr); |
|
memset(&addr, 0, sizeof(addr)); |
|
ssize_t recv_len = socket_recvfrom(sock, data, PACKET_DATA_SIZE, (struct sockaddr*)&addr, &addr_len); |
|
|
|
if (recv_len <= 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "recvfrom failed, error=%zd, sock_err=%d", recv_len, socket_get_error()); |
|
return; |
|
} |
|
|
|
// DUMP: Show received packet content |
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "RECV in:", data, recv_len); |
|
|
|
struct ETCP_DGRAM* pkt = memory_pool_alloc(e_sock->instance->pkt_pool); |
|
if (!pkt) return; |
|
size_t pkt_len=0; |
|
int errorcode=0; |
|
|
|
struct ETCP_LINK* link=etcp_link_find_by_addr(e_sock, &addr); |
|
|
|
// Try normal decryption first if we have an established link with session keys |
|
// This is the common case for data packets and responses |
|
// if (link) { |
|
// link->recv_keepalive = 1; // Link is up after successful initialization - не ставим ап от неизвестных пакетов |
|
// } |
|
|
|
if (link!=NULL && link->etcp!=NULL && link->etcp->crypto_ctx.session_ready) { |
|
sc_status_t dec_rc = sc_decrypt(&link->etcp->crypto_ctx, data, recv_len, (uint8_t*)&pkt->timestamp, &pkt_len); |
|
if (!dec_rc) { |
|
goto process_decrypted; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "etcp: DECRYPT FAIL rc=%d link=%p log=%s sess=%d link_state=%d keepalive=%d enc_errs=%zu my_pub=%016llx peer_pub=%016llx", |
|
dec_rc, link, link->etcp->log_name, link->etcp->crypto_ctx.session_ready, |
|
link->link_state, link->recv_keepalive, link->encrypt_errors, |
|
*(uint64_t*)link->etcp->crypto_ctx.pk->public_key, *(uint64_t*)link->etcp->crypto_ctx.peer_public_key); |
|
} else { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "SKIP normal decrypt: link=%p session_ready=%d — trying init decrypt", |
|
link, link && link->etcp ? link->etcp->crypto_ctx.session_ready : -1); |
|
} |
|
|
|
// Try INIT decryption (for incoming connection requests) |
|
// This handles: no link found, or link without session, or normal decrypt failed |
|
if (recv_len <= SC_PUBKEY_ENC_SIZE + UDP_SC_HDR_SIZE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "packet too small for init, size=%zd from %s", recv_len, sockaddr_storage_to_str(&addr).str); |
|
errorcode=1; |
|
goto ec_fr; |
|
} |
|
|
|
struct secure_channel sc; |
|
sc_init_ctx(&sc, &e_sock->instance->my_keys); |
|
|
|
const uint8_t* salt = data + recv_len - SC_PUBKEY_ENC_SIZE; |
|
const uint8_t* encrypted_pubkey = salt + SC_PUBKEY_ENC_SALT_SIZE; |
|
|
|
uint8_t decrypted_pubkey[SC_PUBKEY_SIZE]; |
|
sc_obfuscate_pubkey(salt, e_sock->instance->my_keys.public_key, encrypted_pubkey, decrypted_pubkey); |
|
|
|
if (sc_set_peer_public_key(&sc, decrypted_pubkey, SC_PEER_PUBKEY_BIN)!=SC_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to set peer public key during init, from %s", sockaddr_storage_to_str(&addr).str); |
|
errorcode=2; |
|
goto ec_fr; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "X25519 decrypt OK from %s", sockaddr_storage_to_str(&addr).str); |
|
if (sc_decrypt(&sc, data, recv_len - SC_PUBKEY_ENC_SIZE, (uint8_t*)&pkt->timestamp, &pkt_len)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to decrypt init packet, from %s", sockaddr_storage_to_str(&addr).str); |
|
errorcode=3; |
|
goto ec_fr; |
|
} |
|
|
|
// INIT decryption succeeded - process packet |
|
if (pkt_len<3) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "too short packet, from %s", sockaddr_storage_to_str(&addr).str); |
|
errorcode=7; |
|
goto ec_fr; |
|
} |
|
if (pkt_len<15) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "decrypted packet too short for INIT header, pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(&addr).str); |
|
errorcode=7; |
|
goto ec_fr; |
|
} |
|
|
|
pkt->data_len=pkt_len-3; |
|
pkt->noencrypt_len=0; |
|
uint8_t code = pkt->data[0]; |
|
uint64_t peer_id = be64toh(*(uint64_t*)(pkt->data + 1)); |
|
if (code == ETCP_PING) { |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "X25519 decrypted: PING from peer=0x%016llx src=%s", |
|
(unsigned long long)peer_id, sockaddr_storage_to_str(&addr).str); |
|
int ret = handle_ping(e_sock, pkt, &addr, decrypted_pubkey, pkt_len); |
|
if (ret) { errorcode = ret; goto ec_fr; } |
|
return; |
|
} |
|
if (code == ETCP_PONG) { |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "X25519 decrypted: PONG from peer=0x%016llx src=%s", |
|
(unsigned long long)peer_id, sockaddr_storage_to_str(&addr).str); |
|
int ret = handle_pong(e_sock, pkt, &addr, pkt_len); |
|
if (ret) { errorcode = ret; goto ec_fr; } |
|
return; |
|
} |
|
if (code!=ETCP_INIT_REQUEST && code!=ETCP_INIT_REQUEST_NOINIT) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "not an init packet: code=0x%02x (expected 0x02/0x04) from %s — packet dropped", |
|
code, sockaddr_storage_to_str(&addr).str); |
|
errorcode=4; |
|
goto ec_fr; |
|
}// не init |
|
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "INIT rejected: socket %s (type=PRIVATE) does not accept incoming connections", e_sock->name); |
|
errorcode=8; |
|
goto ec_fr; |
|
} |
|
if (pkt_len < ETCP_INIT_REQ_SIZE) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "INIT REQUEST too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(&addr).str); |
|
errorcode=7; |
|
goto ec_fr; |
|
} |
|
|
|
// Check allowed keys for incoming connections |
|
struct global_config *global = &e_sock->instance->config->global; |
|
if (!global->allowed_keys_allow_all) { |
|
if (global->allowed_keys_count == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Connection rejected: no allowed_keys configured, allow_all=0"); |
|
errorcode=8; |
|
goto ec_fr; |
|
} |
|
int found = 0; |
|
struct CFG_ALLOWED_KEY *ak = global->allowed_keys; |
|
while (ak) { |
|
if (memcmp(ak->key_bin, sc.peer_public_key, SC_PUBKEY_SIZE) == 0) { |
|
found = 1; |
|
break; |
|
} |
|
ak = ak->next; |
|
} |
|
if (!found) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Connection rejected: peer public key not in allowed_keys list"); |
|
errorcode=8; |
|
goto ec_fr; |
|
} |
|
} |
|
|
|
struct ETCP_INIT_REQUEST_PKT* req = (struct ETCP_INIT_REQUEST_PKT*)pkt->data; |
|
uint32_t session_id = be32toh(*(uint32_t*)req->session_id); |
|
uint16_t mtu = be16toh(*(uint16_t*)req->mtu); |
|
uint16_t src_port = be16toh(*(uint16_t*)req->src_port); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, |
|
"INIT accepted: peer=0x%016llx mtu=%u link=%u sock=%u type=%d local=%d" |
|
" src_ip=%d.%d.%d.%d:%u session=%08x src=%s", |
|
(unsigned long long)peer_id, mtu, req->link_id, req->socket_id, |
|
req->type, req->only_local, |
|
req->src_ipv4[0], req->src_ipv4[1], req->src_ipv4[2], req->src_ipv4[3], |
|
src_port, session_id, sockaddr_storage_to_str(&addr).str); |
|
|
|
struct ETCP_CONN* conn = NULL; |
|
{ |
|
struct ll_entry* e = queue_find_data_by_index(e_sock->instance->connections, (const uint8_t*)&peer_id); |
|
if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; conn = ce->conn; } |
|
} |
|
if (!conn) { |
|
struct ETCP_LINK* ol = etcp_link_find_by_addr(e_sock, &addr); |
|
if (ol && ol->etcp && ol->etcp->peer_node_id == 0) { |
|
conn = ol->etcp; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] reusing unindexed outbound conn for incoming INIT from peer=0x%016llx addr=%s", |
|
conn->log_name, (unsigned long long)peer_id, sockaddr_storage_to_str(&addr).str); |
|
} |
|
} |
|
|
|
int new_conn=0; |
|
if (!conn || conn->peer_node_id!=peer_id) {// создаём новое подключение [new etcp] |
|
new_conn=1; |
|
conn=etcp_connection_create(e_sock->instance,""); |
|
if (!conn) { errorcode=55; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "failed to create connection"); goto ec_fr; } |
|
memcpy(&conn->crypto_ctx, &sc, sizeof(sc)); |
|
etcp_conn_set_peer_node_id(conn, peer_id); |
|
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "New connection received on socket %s: log_name=%s peer_id=%lu peer:%s", e_sock->name, conn->log_name, (unsigned long)peer_id, sockaddr_storage_to_str(&addr).str); |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p total=%d", |
|
ip_to_str(&addr, addr.ss_family).str, peer_id, conn, |
|
queue_entry_count(e_sock->instance->connections)); |
|
} |
|
else {// check keys если существующее подключение |
|
if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) { errorcode=5; DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx", (unsigned long long)peer_id); goto ec_fr; }// коллизия - peer id совпал а ключи разные. |
|
} |
|
|
|
// Check if link already exists (for CHANNEL_INIT recovery) |
|
|
|
struct ETCP_LINK* existing_link = etcp_link_find_by_remote_id(conn, req->link_id); |
|
if (!existing_link) { |
|
existing_link = etcp_link_find_by_addr(e_sock, &addr); |
|
if (existing_link && existing_link->etcp == conn) { |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] found existing outbound link by addr for incoming INIT, reusing link=%p", |
|
conn->log_name, existing_link); |
|
} else { |
|
existing_link = NULL; |
|
} |
|
} |
|
|
|
uint8_t send_reset = 0; |
|
|
|
if (existing_link && existing_link->etcp == conn) {// существующий линк |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] found existing link for id=%d, socket=[%s]", conn->log_name, req->link_id, e_sock->name); |
|
link = existing_link; |
|
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); |
|
if (link->conn) remove_link(link->conn, link->ip_port_hash); // remove old connection from old socket |
|
link->conn=e_sock; |
|
memcpy(&link->remote_addr, &addr, sizeof(addr)); |
|
link->ip_port_hash = sockaddr_hash(&addr); |
|
if (insert_link(link->conn, link) < 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to reinsert link after addr change"); |
|
goto ec_fr; |
|
} |
|
// Cancel pending NAT check and reset status |
|
{ struct TOPO_GROUP* g = topo_groups_get_default(link->etcp->instance->topo_groups); |
|
if (g) route_ping_cancel_for_conn(g, link->etcp); } |
|
link->nat_check_status = NAT_CHECK_NONE; |
|
link->nat_type = NAT_TYPE_UNKNOWN; |
|
} |
|
|
|
// Link exists - reuse it for recovery |
|
link->remote_link_id = req->link_id; |
|
link->remote_socket_id = req->socket_id; |
|
link->remote_only_local = req->only_local; |
|
link->remote_type = req->type; |
|
|
|
// ── Collision handling ── |
|
if (req->collision) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave", |
|
conn->log_name, (unsigned long long)peer_id); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
send_reset = 1; |
|
} else { |
|
// Check if WE have an outbound (master) link on this conn |
|
struct ETCP_LINK* ml = conn->links; |
|
while (ml) { if (ml->is_server == 0) break; ml = ml->next; } |
|
if (ml) { |
|
if (conn->instance->node_id < peer_id) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT", |
|
conn->log_name, (unsigned long long)conn->instance->node_id, (unsigned long long)peer_id); |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
etcp_link_send_init(ml, 0, 1); |
|
return; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: peer smaller, yielding master, processing as slave", |
|
conn->log_name); |
|
} |
|
|
|
// Normal reinit check (skip if reset_done=1 unless session_id changed) |
|
if (!conn->reset_done) { |
|
if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { |
|
send_reset = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT existing link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d", |
|
conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
} |
|
} else if (conn->session_id != session_id) { |
|
send_reset = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT existing link (session changed): sess=%08x→%08x", |
|
conn->log_name, conn->session_id, session_id); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
} else { |
|
send_reset = 0; |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "same session_id, skip reinit"); |
|
} |
|
} |
|
|
|
// Cancel existing timers |
|
if (link->init_timer) { |
|
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer); |
|
link->init_timer = NULL; |
|
} |
|
} else { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] NO existing link for id=%d, socket=[%s]", conn->log_name, req->link_id, e_sock->name); |
|
// Create new link |
|
link = etcp_link_new(conn, e_sock, &addr, 1); |
|
if (!link) { if (new_conn) etcp_connection_close(conn); errorcode=66; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "failed to create link for connection"); goto ec_fr; }// облом |
|
link->remote_link_id = req->link_id; |
|
link->remote_socket_id = req->socket_id; |
|
link->remote_only_local = req->only_local; |
|
link->remote_type = req->type; |
|
|
|
// ── Collision handling ── |
|
if (req->collision) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave", |
|
conn->log_name, (unsigned long long)peer_id); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
send_reset = 1; |
|
} else { |
|
// Check if WE have an outbound (master) link |
|
struct ETCP_LINK* ml = conn->links; |
|
while (ml) { if (ml->is_server == 0) break; ml = ml->next; } |
|
if (ml && conn->instance->node_id < peer_id) { |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] COLLISION: we=master node_id=%016llx < peer=%016llx — sending collision INIT", |
|
conn->log_name, (unsigned long long)conn->instance->node_id, (unsigned long long)peer_id); |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
etcp_link_send_init(ml, 0, 1); |
|
return; |
|
} |
|
|
|
// Normal reinit check (skip if reset_done=1 unless session_id changed) |
|
if (!conn->reset_done) { |
|
if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || !conn->got_initial_pkt) { |
|
send_reset = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT new link: code=0x%02x sess=%08x→%08x got_init=%d initialized=%d links_up=%d", |
|
conn->log_name, code, conn->session_id, session_id, conn->got_initial_pkt, conn->initialized, conn->links_up); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
} |
|
} else if (conn->session_id != session_id) { |
|
send_reset = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[%s] REINIT new link (session changed): sess=%08x→%08x", |
|
conn->log_name, conn->session_id, session_id); |
|
conn->session_id = session_id; |
|
etcp_conn_reinit(conn); |
|
} |
|
} |
|
link->keepalive_interval=(req->keepalive[0]<<8) | req->keepalive[1]; |
|
link->recovery_interval=((req->recovery[0]<<8) | req->recovery[1])*100;// timebase в link, timebase/100 в кодограмме |
|
if (link->keepalive_interval < 10) link->keepalive_interval = 10; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "set keepalive for link=%d", link->keepalive_interval); |
|
} |
|
|
|
link->mtu_remote = (req->mtu[0] << 8) | req->mtu[1]; |
|
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU; |
|
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote; |
|
etcp_update_mtu(link->etcp); |
|
|
|
uint32_t req_src_ip; memcpy(&req_src_ip, req->src_ipv4, 4); |
|
uint16_t req_src_port = be16toh(*(uint16_t*)req->src_port); |
|
|
|
memcpy(conn->peer_ed25519_pubkey, req->ed25519_pubkey, SC_PUBKEY_SIZE); |
|
|
|
send_init_response(e_sock, pkt, link, conn, req, &addr, send_reset, req_src_ip, req_src_port, pkt_len); |
|
return; |
|
|
|
|
|
process_decrypted: |
|
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "Decrypt ok - normal pkt"); |
|
if (pkt_len<3) { errorcode=46; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "decrypted packet too small, size=%zu", pkt_len); goto ec_fr; } |
|
pkt->data_len=pkt_len-3; |
|
pkt->noencrypt_len=0; |
|
pkt->link=link; |
|
|
|
link->remote_keepalive = pkt->flag_up; |
|
|
|
/* restore recv_keepalive BEFORE computing link_status = remote && local */ |
|
if (link->recv_keepalive != 1) { |
|
link->recv_keepalive = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) status changed to UP - packet received", |
|
link->etcp->log_name, link, link->local_link_id); |
|
} |
|
|
|
int was_up = link->link_status; |
|
link->link_status = link->remote_keepalive && link->recv_keepalive; |
|
if (link->link_status && !was_up && link->initialized) |
|
loadbalancer_link_ready(link); |
|
|
|
link->last_recv_local_time=get_time_tb(); |
|
link->last_recv_timestamp=pkt->timestamp; |
|
link->last_recv_updated=1; |
|
|
|
// Count decrypted bytes |
|
link->total_decrypted += pkt->data_len; |
|
|
|
size_t offset = 0; |
|
uint8_t pkt_code = pkt->data[offset++]; |
|
|
|
if (pkt_code == ETCP_KEEPALIVE) { |
|
if (pkt->data_len >= 3) { |
|
uint16_t peer_period = pkt->data[1] | ((uint16_t)pkt->data[2] << 8); |
|
link->keepalive_timeout = (uint32_t)peer_period * KA_TIMEOUT_MULT; |
|
} |
|
link->keepalive_recv_count++; |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return; // KA handled, nothing more to process |
|
} |
|
|
|
if (pkt_code == ETCP_INIT_RESPONSE || pkt_code == ETCP_INIT_RESPONSE_NOINIT) { |
|
int ret = handle_init_response_client(e_sock, pkt, link, pkt_code, pkt_len); |
|
if (ret) { errorcode = ret; goto ec_fr; } |
|
return; |
|
} |
|
|
|
if (link->link_state == 2) {// из recovery получен нормальный пакет - восстанавливаем линк в нормальный режим |
|
start_keepalive_timer(link); |
|
etcp_link_send_keepalive(link); // Start keepalive timer |
|
link->link_state = 3; // connected |
|
} |
|
|
|
|
|
// log_dump("RECV decrypted:", pkt->data, pkt->data_len, link); |
|
|
|
if (link->link_state == 3) { |
|
if (memory_pool_is_freed(e_sock->instance->pkt_pool, pkt)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt=%p ALREADY FREED in pkt_pool — HALTING", (void*)pkt); |
|
volatile int _halt = 1; while (_halt) {} |
|
} |
|
etcp_conn_input(pkt); |
|
} else memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return; |
|
|
|
ec_fr: |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "error %d, from %s", errorcode, sockaddr_storage_to_str(&addr).str); |
|
e_sock->pkt_format_errors++; |
|
e_sock->errorcode=errorcode; |
|
memory_pool_free(e_sock->instance->pkt_pool, pkt); |
|
return; |
|
} |
|
|
|
// Initialize only sockets (servers for incoming connections) |
|
// Called before topo_group_init() to populate etcp_sockets for nodeinfo |
|
// Returns: 0 = all OK, 1 = partial success (some sockets failed), -1 = fatal error (no sockets) |
|
int init_sockets(struct UTUN_INSTANCE* instance) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!instance || !instance->config) return -1; |
|
if (instance->etcp_sockets || instance->stcp_server) { |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets already initialized, skipping"); |
|
return 0; |
|
} |
|
|
|
struct utun_config* config = instance->config; |
|
int success_count = 0; |
|
int fail_count = 0; |
|
|
|
// Create sockets for servers (incoming connections) |
|
struct CFG_SERVER* server = config->servers; |
|
while (server) { |
|
// Auto-detect local IP for public servers with 0.0.0.0 |
|
uint32_t default_ip = 0; |
|
uint8_t default_ip6[16] = {0}; |
|
int have_default_ip6 = 0; |
|
if (server->type == CFG_SERVER_TYPE_PUBLIC) { |
|
if (server->ip.ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)&server->ip; |
|
if (sin->sin_addr.s_addr == 0) { |
|
struct sockaddr_storage remote; |
|
memset(&remote, 0, sizeof(remote)); |
|
struct sockaddr_in* rem_sin = (struct sockaddr_in*)&remote; |
|
rem_sin->sin_family = AF_INET; |
|
rem_sin->sin_port = htons(53); |
|
inet_pton(AF_INET, "8.8.8.8", &rem_sin->sin_addr); |
|
struct sockaddr_storage local; |
|
if (get_outgoing_local_ip(&server->ip, &remote, &local) == 0) { |
|
struct sockaddr_in* local_sin = (struct sockaddr_in*)&local; |
|
default_ip = local_sin->sin_addr.s_addr; |
|
} |
|
if (default_ip == 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "Failed to detect default route IP for server %s", server->name); |
|
} |
|
} |
|
} else if (server->ip.ss_family == AF_INET6) { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&server->ip; |
|
if (memcmp(&sin6->sin6_addr, &in6addr_any, 16) == 0) { |
|
struct sockaddr_storage remote; |
|
memset(&remote, 0, sizeof(remote)); |
|
struct sockaddr_in6* rem_sin6 = (struct sockaddr_in6*)&remote; |
|
rem_sin6->sin6_family = AF_INET6; |
|
rem_sin6->sin6_port = htons(53); |
|
inet_pton(AF_INET6, "2001:4860:4860::8888", &rem_sin6->sin6_addr); |
|
struct sockaddr_storage local; |
|
if (get_outgoing_local_ip(&server->ip, &remote, &local) == 0) { |
|
struct sockaddr_in6* local_sin6 = (struct sockaddr_in6*)&local; |
|
memcpy(default_ip6, &local_sin6->sin6_addr, 16); |
|
have_default_ip6 = 1; |
|
} |
|
if (!have_default_ip6) { |
|
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "Failed to detect default route IPv6 for server %s", server->name); |
|
} |
|
} |
|
} |
|
} |
|
|
|
// TCP transport: create stcp_server instead of UDP socket |
|
if (server->transport) { |
|
uint16_t port = ntohs(((struct sockaddr_in*)&server->ip)->sin_port); |
|
struct stcp_link_config scfg = {.ua = instance->ua, .my_keys = &instance->my_keys, .inst = instance}; |
|
struct stcp_server *tsrv = stcp_server_listen(&scfg, port, tcp_server_on_link, instance); |
|
if (!tsrv) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create TCP server for %s", server->name); |
|
fail_count++; |
|
} else { |
|
instance->stcp_server = tsrv; |
|
success_count++; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server %s on port %u", server->name, port); |
|
} |
|
server = server->next; |
|
continue; |
|
} |
|
|
|
struct ETCP_SOCKET* e_sock = etcp_socket_add(instance, server); |
|
if (e_sock && default_ip != 0) { |
|
e_sock->local_defaultroute_ip = default_ip; |
|
if (server->ip.ss_family == AF_INET) { |
|
struct in_addr addr; |
|
addr.s_addr = default_ip; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Server %s type %d ip=%s", server->name, server->type, ip_to_str(&addr, AF_INET).str); |
|
} |
|
} |
|
if (e_sock && have_default_ip6) { |
|
memcpy(e_sock->local_defaultroute_ip6, default_ip6, 16); |
|
} |
|
if (!e_sock) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create socket for server %s", server->name); |
|
fail_count++; |
|
server = server->next; |
|
continue; |
|
} |
|
success_count++; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized server %s on %s (links: %zu)", |
|
server->name, sockaddr_storage_to_str(&server->ip).str, e_sock->num_channels); |
|
server = server->next; |
|
} |
|
|
|
if (success_count == 0 && fail_count > 0) return -1; // All failed - fatal |
|
if (fail_count > 0) return 1; // Partial failure - non-fatal warning |
|
return 0; // All OK |
|
} |
|
|
|
int init_connections(struct UTUN_INSTANCE* instance) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); |
|
if (!instance || !instance->config) return -1; |
|
|
|
struct utun_config* config = instance->config; |
|
int socket_result = 0; |
|
|
|
// If sockets already exist (created by init_sockets), check stored status |
|
if (instance->etcp_sockets || instance->stcp_server) { |
|
if (instance->socket_init_status == 1) { |
|
socket_result = 1; |
|
} |
|
} else { |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets not initialized, calling init_sockets()"); |
|
socket_result = init_sockets(instance); |
|
instance->socket_init_status = socket_result; |
|
if (socket_result < 0) { |
|
return -1; |
|
} |
|
} |
|
|
|
// Initialize clients - create outgoing connections |
|
struct CFG_CLIENT* client = config->clients; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, total_conns=%d", |
|
instance, config, config ? config->clients : NULL, instance ? queue_entry_count(instance->connections) : -1); |
|
while (client) { |
|
// Check if client has required configuration |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Client %s - keepalive=%d, links=%p, peer_key_len=%zu", |
|
client->name, client->keepalive, client->links, |
|
strlen(client->peer_public_key_hex)); |
|
|
|
// Create ETCP connection for this client |
|
struct ETCP_CONN* etcp_conn = etcp_connection_create(instance, client->name); |
|
|
|
if (!etcp_conn) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "Failed to create ETCP connection for client %s", client->name); |
|
client = client->next; |
|
continue; |
|
} |
|
|
|
// Generate session_id for this client connection |
|
if (random_bytes((uint8_t*)&etcp_conn->session_id, sizeof(etcp_conn->session_id)) != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "Failed to generate session_id for client %s", client->name); |
|
etcp_connection_close(etcp_conn); |
|
client = client->next; |
|
continue; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Client %s session_id=%08x", client->name, etcp_conn->session_id); |
|
|
|
// Initialize crypto context for this connection |
|
if (sc_init_ctx(&etcp_conn->crypto_ctx, &instance->my_keys) != SC_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to initialize crypto context for client %s", client->name); |
|
etcp_connection_close(etcp_conn); |
|
client = client->next; |
|
continue; |
|
} |
|
// If client has peer public key configured, set it |
|
if (strlen(client->peer_public_key_hex) > 0) { |
|
// For now, set peer node ID to indicate we have peer key |
|
// The actual peer key will be exchanged during connection establishment |
|
etcp_conn->peer_node_id = 1; // Simple indicator |
|
etcp_update_log_name(etcp_conn); // Update log_name with peer_node_id |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "setting peer public key for client %s", client->name); |
|
// Set peer public key (assuming hex format) |
|
if (sc_set_peer_public_key(&etcp_conn->crypto_ctx, (const uint8_t*)client->peer_public_key_hex, 1) != SC_OK) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to set peer public key for client %s", client->name); |
|
} else { |
|
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "successfully set peer public key for client %s", client->name); |
|
} |
|
} else { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "no peer public key configured for client %s", client->name); |
|
} |
|
|
|
etcp_set_routing_exchange_state(etcp_conn, 1); // инициируем обмен маршрутами |
|
|
|
// Create links for this client |
|
struct CFG_CLIENT_LINK* client_link = client->links; |
|
while (client_link) { |
|
// Find the local server for this link |
|
struct CFG_SERVER* local_server = client_link->local_srv; |
|
if (!local_server) { |
|
client_link = client_link->next; |
|
continue; |
|
} |
|
|
|
// Find the socket for this server |
|
struct ETCP_SOCKET* e_sock = NULL; |
|
struct ETCP_SOCKET* sock = instance->etcp_sockets; |
|
while (sock) { |
|
if (sock->local_addr.ss_family == local_server->ip.ss_family) { |
|
if (sock->local_addr.ss_family == AF_INET) { |
|
struct sockaddr_in* sock_addr = (struct sockaddr_in*)&sock->local_addr; |
|
struct sockaddr_in* srv_addr = (struct sockaddr_in*)&local_server->ip; |
|
if (sock_addr->sin_addr.s_addr == srv_addr->sin_addr.s_addr && |
|
sock_addr->sin_port == srv_addr->sin_port) { |
|
e_sock = sock; |
|
break; |
|
} |
|
} else if (sock->local_addr.ss_family == AF_INET6) { |
|
struct sockaddr_in6* sock_addr6 = (struct sockaddr_in6*)&sock->local_addr; |
|
struct sockaddr_in6* srv_addr6 = (struct sockaddr_in6*)&local_server->ip; |
|
if (memcmp(&sock_addr6->sin6_addr, &srv_addr6->sin6_addr, 16) == 0 && |
|
sock_addr6->sin6_port == srv_addr6->sin6_port) { |
|
e_sock = sock; |
|
break; |
|
} |
|
} |
|
} |
|
sock = sock->next; |
|
} |
|
|
|
if (local_server->transport) { |
|
if (strlen(client->peer_public_key_hex) == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "TCP client %s has no peer key", client->name); |
|
client_link = client_link->next; continue; |
|
} |
|
// TCP transport — use stcp_link instead of etcp_link_new |
|
uint16_t rport = ntohs(((struct sockaddr_in*)&client_link->remote_addr)->sin_port); |
|
struct stcp_link_config scfg = { |
|
.ua = instance->ua, .my_keys = &instance->my_keys, .inst = instance, |
|
.peer_pubkey = (const uint8_t*)client->peer_public_key_hex, |
|
.peer_pubkey_mode = 1, // hex |
|
.remote_addr = &client_link->remote_addr, .remote_port = rport |
|
}; |
|
struct stcp_link *slink = stcp_link_connect(&scfg); |
|
if (!slink) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create TCP link for client %s", client->name); |
|
client_link = client_link->next; continue; |
|
} |
|
etcp_conn->transport_link = slink; |
|
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP link created for client %s", client->name); |
|
client_link = client_link->next; continue; |
|
} |
|
|
|
if (!e_sock) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "No socket found for client %s link", client->name); |
|
client_link = client_link->next; |
|
continue; |
|
} |
|
|
|
// Create link for this client connection |
|
struct ETCP_LINK* link = etcp_link_new(etcp_conn, e_sock, &client_link->remote_addr, 0); // 0 = client initiates |
|
if (!link) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create link for client %s", client->name); |
|
client_link = client_link->next; |
|
continue; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Created link %p for client %s, socket=%p", |
|
link, client->name, e_sock); |
|
|
|
client_link = client_link->next; |
|
} |
|
|
|
client = client->next; |
|
} |
|
|
|
// If there are clients configured but no connections created, that's an error |
|
// If there are no clients (server-only mode), 0 connections is OK (server will accept incoming) |
|
int total_conns = queue_entry_count(instance->connections); |
|
if (total_conns == 0 && config->clients != NULL) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Clients configured but no connections initialized"); |
|
return -1; |
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized %d connections", total_conns); |
|
// Return 1 if there was a partial socket initialization error |
|
if (socket_result == 1) return 1; |
|
return 0; |
|
}
|
|
|