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.
 
 
 
 
 
 

2842 lines
145 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 "etcp.h"
#include "stcp_link.h"
#include "stcp.h"
#include "stcp_client.h"
#include "socks_client.h"
#include "topo_node.h"
#include "topo_group.h"
#include "node_conn_direct.h"
#include "nat_detection.h"
#include "../lib/memory_pool.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "etcp_loadbalancer.h"
#include "etcp_keepalive.h"
#ifdef UTUN_HAVE_STANDBY
#include "standby.h"
#endif
#include <stdlib.h>
#include <time.h>
#include "../lib/mem.h"
#include "etcp.h"
// Читаемое имя типа сокета (CFG_SERVER_TYPE_*) для логов.
static const char* server_type_str(uint8_t type) {
static const char* names[] = {"UNKNOWN", "PUBLIC", "NAT", "PRIVATE", "LOCAL"};
return type < 5 ? names[type] : "?";
}
// Читаемое имя NAT-типа (NAT_TYPE_*/NAT_VERIFIED_*) для логов.
static const char* nat_type_str(uint8_t nat_type) {
if (nat_type >= 4 && nat_type <= 7) nat_type -= 4;
static const char* names[] = {"UNKNOWN", "EIM", "STRICT", "DIRECT"};
return nat_type < 4 ? names[nat_type] : "?";
}
/* ── Forward declarations ────────────────────────────────────────────
* Только static-функции, используемые раньше своего определения.
* Не-static функции уже объявлены в etcp_connections.h. */
static void tcp_link_close_cb(struct stcp_link *sl, int err, void *arg);
static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision);
static void etcp_link_init_timer_cbk(void* arg);
static void burst_resp_timeout_cb(void* arg);
static int etcp_tcp_send(struct ETCP_DGRAM* dgram);
static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data, ssize_t recv_len, struct sockaddr_storage addr);
static void socks_etcp_read_callback(socket_t sock, void* arg);
static void socks_relay_ready_cb(struct socks_udp* s, int err, void* arg);
// Коллбэк готовности SOCKS5 UDP ASSOCIATE: только логирование (sendto сам ждёт ready).
static void socks_relay_ready_cb(struct socks_udp* s, int err, void* arg) {
struct ETCP_SOCKET* e_sock = (struct ETCP_SOCKET*)arg;
(void)s;
if (err != 0) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks UDP associate error for socket %s err=%d", e_sock->name, err); return; }
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socks UDP relay ready for socket %s", e_sock->name);
}
// Чтение с локального UDP-сокета SOCKS5 UDP ASSOCIATE: снять SOCKS-заголовок → реальный src,
// дальше обычный разбор через etcp_process_packet.
static void socks_etcp_read_callback(socket_t sock, void* arg) {
struct ETCP_SOCKET* e_sock = (struct ETCP_SOCKET*)arg;
if (!e_sock || !e_sock->socks_udp) return;
uint8_t raw[PACKET_DATA_SIZE + 64];
uint8_t payload[PACKET_DATA_SIZE];
struct sockaddr_storage relay;
socklen_t rl = sizeof(relay);
memset(&relay, 0, sizeof(relay));
ssize_t raw_len = socket_recvfrom(sock, raw, sizeof(raw), (struct sockaddr*)&relay, &rl);
if (raw_len <= 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "socks recvfrom failed, error=%zd, sock_err=%d", raw_len, socket_get_error());
return;
}
struct sockaddr_storage src;
ssize_t plen = socks_udp_unwrap(e_sock->socks_udp, raw, (size_t)raw_len, payload, sizeof(payload), &src);
if (plen <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "socks unwrap failed len=%zd from %s", plen, sockaddr_storage_to_str(&relay).str);
return;
}
etcp_process_packet(e_sock, payload, plen, src);
}
// STCP-сервер: новое входящее TCP-соединение → найти/создать ETCP_CONN по node_id
// pubkey пира, создать TCP-линк и поднять его (ready).
void tcp_server_on_link(struct stcp_link *link, struct ETCP_SOCKET *tcp_sock) {
struct UTUN_INSTANCE *inst = tcp_sock->instance;
const uint8_t *pubkey = stcp_link_get_peer_pubkey(link);
if (!pubkey) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: no peer pubkey"); return; }
uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: accepted TCP from peer=0x%016llx", (unsigned long long)node_id);
struct ETCP_CONN *conn = instance_find_conn(inst, node_id);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: %s conn=%p peer=0x%016llx pubkey=%016llx",
conn ? "existing" : "NEW", (void*)conn, (unsigned long long)node_id, *(const uint64_t*)pubkey);
if (!conn) {
conn = etcp_connection_create(inst, NULL);
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: etcp_connection_create failed"); return; }
conn->peer_node_id = node_id;
if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: sc_init_ctx failed");
etcp_connection_close(conn); return;
}
if (sc_set_peer_public_key(&conn->crypto_ctx, pubkey, SC_PEER_PUBKEY_BIN) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: sc_set_peer_public_key failed");
etcp_connection_close(conn); return;
}
etcp_update_log_name(conn);
{ struct conn_queue_entry* ce = (struct conn_queue_entry*)conn->conn_queue_entry->data; ce->peer_node_id = node_id; queue_remove_data(conn->conn_queue, conn->conn_queue_entry); queue_data_put_with_index(conn->conn_queue, conn->conn_queue_entry); }
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: new ETCP_CONN peer=0x%016llx", (unsigned long long)node_id);
}
const struct sockaddr_storage* ra = stcp_link_get_peer_addr(link);
if (!ra) ra = stcp_link_get_remote_addr(link);
struct ETCP_LINK *tlink = etcp_link_new(conn, tcp_sock, ra, 1);
if (!tlink) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_server_on_link: etcp_link_new failed"); return; }
tlink->tcp_link = link;
tlink->is_server = 1;
stcp_link_set_etcp_conn(link, conn);
stcp_link_set_etcp_link(link, tlink);
stcp_link_set_on_close(link, tcp_link_close_cb, tlink);
etcp_link_enter_ready_tcp(tlink);
}
// === Burst sender: замер пропускной способности линка ===
// Старт burst: отправка пачки пакетов подряд для измерения bandwidth.
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);
}
// Возможность стартовать burst: линк насыщен (inflight достиг лимита) и прошёл интервал.
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);
}
// Завершение burst: сброс флага, запуск таймера ожидания ответа с замерами пира.
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);
}
// Таймаут ожидания ответа на burst — просто сбрасываем таймер.
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
// Клиент: собрать и отправить INIT_REQUEST (handshake) — заголовок, padding для обфускации
// размера и обфусцированный pubkey. reset=1 → запрос сброса эпохи, collision=1 → ответ на коллизию.
static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t collision) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "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;
if (link->is_tcp) return; // TCP uses STCP handshake instead of ETCP INIT
struct ETCP_DGRAM* dgram = u_malloc(PACKET_DATA_SIZE);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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);
*(uint64_t*)req->reset_id = htobe64(link->etcp->reset_id);
*(uint16_t*)req->mtu = htobe16(link->mtu_local);
*(uint16_t*)req->keepalive = htobe16(link->etcp->instance->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;
// SOCKS5-прокси: пир не может дотянуться до нашего interface_addr — не рекламируем src.
if (link->conn && link->conn->socks_udp) {
memset(req->src_ipv4, 0, 4);
memset(req->src_port, 0, 2);
} else 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 if (link->conn && link->conn->interface_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&link->conn->interface_addr;
memset(req->src_ipv4, 0, 4);
*(uint16_t*)req->src_port = sin6->sin6_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);
req->device_type = link->etcp->instance->client_type;
size_t offset = ETCP_INIT_REQ_SIZE;
// padding
int wire_ov_init = WIRE_OVERHEAD(link->remote_addr.ss_family);
int pad_range = (int)link->handshake_maxsize - (int)link->handshake_minsize;
if (pad_range <= 0) pad_range = 1;
int s = rand() % pad_range + link->handshake_minsize;
int s_max = (int)(link->mtu) - wire_ov_init;
if (s > s_max) s = s_max;
if (s < 0) s = 0;
int to_add = s - offset - wire_ov_init;
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));
#pragma GCC diagnostic push
#pragma GCC diagnostic ignored "-Wstringop-overflow"
memcpy(dgram->data + offset, salt, SC_PUBKEY_ENC_SALT_SIZE);
offset += SC_PUBKEY_ENC_SALT_SIZE;
#pragma GCC diagnostic pop
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);
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "INIT_SEND: salt=%02x%02x%02x%02x%02x%02x%02x%02x obf_pub=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x my_pub=%02x%02x%02x%02x seskey=%02x%02x%02x%02x",
salt[0], salt[1], salt[2], salt[3], salt[4], salt[5], salt[6], salt[7],
obfuscated_pubkey[0], obfuscated_pubkey[1], obfuscated_pubkey[2], obfuscated_pubkey[3],
link->etcp->crypto_ctx.peer_public_key[0], link->etcp->crypto_ctx.peer_public_key[1],
link->etcp->crypto_ctx.peer_public_key[2], link->etcp->crypto_ctx.peer_public_key[3],
link->etcp->instance->my_keys.public_key[0], link->etcp->instance->my_keys.public_key[1],
link->etcp->instance->my_keys.public_key[2], link->etcp->instance->my_keys.public_key[3],
link->etcp->crypto_ctx.session_key[0], link->etcp->crypto_ctx.session_key[1],
link->etcp->crypto_ctx.session_key[2], link->etcp->crypto_ctx.session_key[3]);
#pragma GCC diagnostic push
#pragma GCC diagnostic ignored "-Wstringop-overflow"
memcpy(dgram->data + offset, obfuscated_pubkey, SC_PUBKEY_SIZE);
offset += SC_PUBKEY_SIZE;
#pragma GCC diagnostic pop
dgram->data_len = offset;
if (link->init_retry_count == 0)
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] INIT sent to %s (link=%d, retry=%d, reset=%d)", link->etcp->log_name, sockaddr_storage_to_str(&link->remote_addr).str, link->local_link_id, link->init_retry_count, reset);
else
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] INIT sent to %s (link=%d, retry=%d, reset=%d)", link->etcp->log_name, sockaddr_storage_to_str(&link->remote_addr).str, link->local_link_id, link->init_retry_count, reset);
etcp_encrypt_send(dgram);
u_free(dgram);
link->init_retry_count++;
}
// Периодический повтор INIT до установления связи (с экспоненциальным backoff до INIT_TIMEOUT_MAX).
static void etcp_link_init_timer_cbk(void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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 == LINK_STATE_CONNECTED && link->initialized) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED state=%d init=%d link_status=%d recv_ka=%d remote_ka=%d links_up=%d", link->etcp->log_name, link->link_state, link->initialized, link->link_status, link->recv_keepalive, link->remote_keepalive, link->etcp->links_up); return; }
if (link->etcp->links_up > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED (links_up=%d > 0)", link->etcp->log_name, link->etcp->links_up); return; }
if (link->etcp->fin_wait) { DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] init_timer: SUPPRESSED (fin_wait)", link->etcp->log_name); return; }
if (link->is_tcp) { etcp_tcp_link_start_reconnect(link); return; } // TCP uses STCP handshake, not ETCP INIT
if (link->etcp->got_initial_pkt == 0) etcp_link_send_init(link,1,0);
else etcp_link_send_init(link,0,0);
}
// Сброс таймера INIT на начальный таймаут (используется при каждом новом handshake).
void etcp_link_restart_init_timer(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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");
}
// Клиент: переход в handshake — отправить INIT и запустить таймер повторов.
void etcp_link_enter_init(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!link) return;
int old_state = link->link_state;
link->link_state = LINK_STATE_HANDSHAKE; // handshake
etcp_fire_link_status_cbk(link, old_state, link->link_status);
if (link->is_server != 0) return;
int reset = (link->etcp->got_initial_pkt == 0);
etcp_link_send_init(link, reset, 0);
etcp_link_restart_init_timer(link);
}
// Клиент: повторный handshake после разрыва — уведомить о падении, отправить INIT, снять keepalive.
void etcp_link_enter_reinit(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!link) return;
int old_state = link->link_state;
link->link_state = LINK_STATE_TRY_RECONNECT; // reconnect
etcp_fire_link_status_cbk(link, old_state, link->link_status);
etcp_on_link_down(link->etcp, link);
if (link->is_server != 0) return;
int reset = (link->etcp->got_initial_pkt == 0);
etcp_link_send_init(link, reset, 0);
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);
}
// Собирает 19-байтовый ключ из sockaddr (port+addr+family) для hash-индекса links_queue.
static void sockaddr_to_key(struct sockaddr_storage* addr, uint8_t key[LINK_ADDR_KEY_SIZE], int is_tcp) {
memset(key, 0, LINK_ADDR_KEY_SIZE);
if (!addr) return;
if (addr->ss_family == AF_INET) {
struct sockaddr_in* sa = (struct sockaddr_in*)addr;
memcpy(key, &sa->sin_port, 2);
memcpy(key + 2, &sa->sin_addr.s_addr, 4);
key[18] = 4 | (is_tcp ? 0x80 : 0);
} else {
struct sockaddr_in6* sa = (struct sockaddr_in6*)addr;
memcpy(key, &sa->sin6_port, 2);
memcpy(key + 2, &sa->sin6_addr, 16);
key[18] = 6 | (is_tcp ? 0x80 : 0);
}
}
// find_link_index, realloc_links, insert_link, remove_link — УДАЛЕНЫ. Заменены на ll_queue с хеш-индексом.
// Вставка линка в links_queue сокета по hash-ключу адреса; при дубле адреса другого коннекта — коллизия.
static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!e_sock || !link || !e_sock->links_queue) return -1;
uint8_t key[LINK_ADDR_KEY_SIZE];
sockaddr_to_key(&link->remote_addr, key, link->is_tcp);
struct ll_entry* dup_qe = queue_find_data_by_index(e_sock->links_queue, key);
if (dup_qe) {
struct link_queue_entry* dup_lqe = (struct link_queue_entry*)dup_qe->data;
if (dup_lqe->link && dup_lqe->link->etcp != link->etcp) {
struct ETCP_CONN* new_c = link->etcp;
struct ETCP_CONN* old_c = dup_lqe->link->etcp;
uint64_t new_qkey = 0, old_qkey = 0;
if (new_c->conn_queue_entry) new_qkey = ((struct conn_queue_entry*)new_c->conn_queue_entry->data)->peer_node_id;
if (old_c->conn_queue_entry) old_qkey = ((struct conn_queue_entry*)old_c->conn_queue_entry->data)->peer_node_id;
DEBUG_ERROR(DEBUG_CATEGORY_ETCP,
"!!!!!!!!!!!! LINK ADDR COLLISION !!!!!!!!!!!! addr=%s "
"new: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx "
"old: conn=%p [%s] peer=0x%016llx state=%d init=%d up=%d srv=%d tcp=%d qkey=0x%016llx",
sockaddr_storage_to_str(&link->remote_addr).str,
(void*)new_c, new_c->log_name, (unsigned long long)new_c->peer_node_id,
new_c->state, new_c->initialized, new_c->links_up, link->is_server, link->is_tcp,
(unsigned long long)new_qkey,
(void*)old_c, old_c->log_name, (unsigned long long)old_c->peer_node_id,
old_c->state, old_c->initialized, old_c->links_up,
dup_lqe->link->is_server, dup_lqe->link->is_tcp,
(unsigned long long)old_qkey);
return -1;
}
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "insert_link_queue: replacing stale DUP addr in [%s] old_link=%p", e_sock->name, dup_lqe->link);
if (dup_lqe->link) dup_lqe->link->link_queue_entry = NULL;
queue_remove_data(e_sock->links_queue, dup_qe);
queue_entry_free(dup_qe);
}
struct ll_entry* qe = queue_entry_new(sizeof(struct link_queue_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "queue_entry_new failed"); return -1; }
struct link_queue_entry* lqe = (struct link_queue_entry*)qe->data;
memcpy(lqe->key, key, LINK_ADDR_KEY_SIZE);
lqe->link = link;
link->link_queue_entry = qe;
queue_data_put_with_index(e_sock->links_queue, qe);
return 0;
}
// Удаление линка из links_queue сокета (освобождение hash-записи).
static void remove_link_from_queue(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!link || !link->link_queue_entry) return;
if (!link->conn || !link->conn->links_queue) return;
queue_remove_data(link->conn->links_queue, link->link_queue_entry);
queue_entry_free(link->link_queue_entry);
link->link_queue_entry = NULL;
}
// Сравнение двух sockaddr_storage (family + addr + port) без учёта прочих полей.
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;
}
// Поиск линка в сокете по remote-адресу (через hash-индекс links_queue).
struct ETCP_LINK* etcp_link_find_by_addr(struct ETCP_SOCKET* e_sock, struct sockaddr_storage* addr, int is_tcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!e_sock || !addr || !e_sock->links_queue) return NULL;
uint8_t key[LINK_ADDR_KEY_SIZE];
sockaddr_to_key(addr, key, is_tcp);
struct ll_entry* e = queue_find_data_by_index(e_sock->links_queue, key);
if (!e) return NULL;
struct link_queue_entry* lqe = (struct link_queue_entry*)e->data;
return lqe->link;
}
/* Перепривязывает UDP-линк на новый remote-адрес (пир сменил NAT-адрес, линк ещё down).
* Пере-индексирует линк в socket->links_queue под новым ключом. */
static void etcp_link_update_remote_addr(struct ETCP_LINK* link, struct ETCP_SOCKET* e_sock,
const struct sockaddr_storage* new_addr) {
if (!link || !e_sock || !new_addr) return;
if (link->is_tcp) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] update_remote_addr: TCP link skip", link->etcp->log_name);
return;
}
if (link->conn != e_sock) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] update_remote_addr: link socket [%s] != recv socket [%s], skip",
link->etcp->log_name, link->conn ? link->conn->name : "-", e_sock->name);
return;
}
remove_link_from_queue(link);
memcpy(&link->remote_addr, new_addr, sizeof(struct sockaddr_storage));
if (insert_link_queue(e_sock, link) < 0)
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] update_remote_addr: re-index FAILED at %s",
link->etcp->log_name, sockaddr_storage_to_str(new_addr).str);
}
/* Выбор линка для отправки collision INIT: outbound-линк (is_server==0), чей remote_addr
* совпадает с адресом пришедшего INIT (симметричный сокет). NULL — симметричного сокета
* нет (асимметрия/NAT) → коллизия не срабатывает, inbound принимается как есть. */
struct ETCP_LINK* etcp_select_collision_link(struct ETCP_CONN* conn, const struct sockaddr_storage* addr) {
if (!conn || !addr) return NULL;
struct ETCP_LINK* l = conn->links;
while (l) {
if (l->is_server == 0 && sockaddr_equal(&l->remote_addr, addr)) return l;
l = l->next;
}
return NULL;
}
// Поиск линка по remote_link_id (id, под которым линк известен пиру).
struct ETCP_LINK* etcp_link_find_by_remote_id(struct ETCP_CONN* conn, uint8_t remote_link_id) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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;
}
// Выделение свободного local_link_id (0-255): битовая маска занятых id + циклический поиск.
int etcp_find_free_local_link_id(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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, начиная с last_assigned_link_id + 1 (циклический перебор)
for (int i = 0; i < 256; i++) {
int id = ((etcp->last_assigned_link_id + 1 + i) & 0xFF);
if (id == 0) continue;
if (!(used_ids[id >> 3] & (1 << (id & 7)))) {
etcp->last_assigned_link_id = (uint8_t)id;
return id;
}
}
// Все id заняты
return -1;
}
// Создание UDP-сокета приёма/отправки кодограмм: bind к interface_addr, настройка буферов,
// регистрация в uasync и добавление в instance->etcp_sockets.
struct ETCP_SOCKET* etcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_SERVER* server) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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;
uint8_t socks_enabled = server->socks_enabled;
char* name = server->name;
struct ETCP_SOCKET* e_sock = u_calloc(1, sizeof(struct ETCP_SOCKET));
if (!e_sock) {
DEBUG_ERROR(DEBUG_CATEGORY_SYS, "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_SOCKET, "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_addr_nl(netif_index, !temp, v6addr) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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;
e_sock->netif_index = netif_index;
if (family != AF_INET && family != AF_INET6) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "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_SOCKET, "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_SOCKET, "Failed to set non-blocking mode");
}
if (family == AF_INET6) {
int v6only = 1;
if (setsockopt(e_sock->fd, IPPROTO_IPV6, IPV6_V6ONLY, &v6only, sizeof(v6only)) < 0) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "Failed to set IPV6_V6ONLY: %s", socket_strerror(socket_get_error()));
}
}
// 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_SOCKET, "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
// Определяем interface_addr до bind, чтобы забиндить на конкретный адрес
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) {
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_SOCKET, "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) {
uint16_t v6if = (netif_index > 0) ? netif_index : get_default_route_netif_index(AF_INET6);
if (v6if > 0) {
uint8_t v6addr[16];
int got = get_interface_ipv6_addr_nl(v6if, 1, v6addr);
if (got != 0) got = get_interface_ipv6_addr_nl(v6if, 0, v6addr);
if (got == 0) {
struct sockaddr_in6* if6 = (struct sockaddr_in6*)&e_sock->interface_addr;
if6->sin6_family = AF_INET6;
memcpy(&if6->sin6_addr, v6addr, 16);
if6->sin6_port = sin6->sin6_port;
}
}
if (e_sock->interface_addr.ss_family == 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);
if (get_outgoing_local_ip(ip, &remote, &e_sock->interface_addr) == 0) {
struct sockaddr_in6* if6 = (struct sockaddr_in6*)&e_sock->interface_addr;
if6->sin6_port = sin6->sin6_port;
}
}
}
}
}
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;
}
}
// Bind — используем interface_addr если определён, иначе ip из конфига
{
struct sockaddr_storage* bind_addr = (e_sock->interface_addr.ss_family != 0) ? &e_sock->interface_addr : ip;
if (bind_addr && bind_addr->ss_family != 0) {
memcpy(&e_sock->local_addr, bind_addr, sizeof(struct sockaddr_storage));
socklen_t addr_len = (bind_addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6);
if (bind(e_sock->fd, (struct sockaddr*)bind_addr, addr_len) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[ETCP] Failed to bind socket to address family %d: %s",
bind_addr->ss_family, socket_strerror(socket_get_error()));
if (bind_addr->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)bind_addr;
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 (bind_addr->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)bind_addr;
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 (bind_addr->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)bind_addr;
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));
} else if (bind_addr->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)bind_addr;
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));
}
}
}
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_ETCP, "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;
e_sock->links_queue = queue_new(instance->ua, 256, offsetof(struct link_queue_entry, key), LINK_ADDR_KEY_SIZE, "links");
if (!e_sock->links_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "Failed to create links_queue for socket %s", e_sock->name);
socket_close_wrapper(e_sock->fd);
u_free(e_sock);
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "Add Socket type=%s", server_type_str(type));
e_sock->next = instance->etcp_sockets;
instance->etcp_sockets = e_sock;
if (socks_enabled) {
// SOCKS5 UDP ASSOCIATE: весь сокет туннелируется через прокси.
struct global_config* gc = &instance->config->global;
if (!gc->socks_host[0] || !gc->socks_port) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socket %s socks=yes but [socks] addr not configured", e_sock->name);
socket_close_wrapper(e_sock->fd);
queue_free(e_sock->links_queue);
instance->etcp_sockets = e_sock->next;
u_free(e_sock);
return NULL;
}
struct socks_cfg scfg;
memset(&scfg, 0, sizeof(scfg));
strncpy(scfg.host, gc->socks_host, sizeof(scfg.host) - 1);
scfg.port = gc->socks_port;
strncpy(scfg.user, gc->socks_username, sizeof(scfg.user) - 1);
strncpy(scfg.pass, gc->socks_password, sizeof(scfg.pass) - 1);
e_sock->socks_udp = socks_udp_associate(instance->ua, &scfg, e_sock->fd, socks_relay_ready_cb, e_sock);
if (!e_sock->socks_udp) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "socks_udp_associate failed for socket %s", e_sock->name);
socket_close_wrapper(e_sock->fd);
queue_free(e_sock->links_queue);
instance->etcp_sockets = e_sock->next;
u_free(e_sock);
return NULL;
}
e_sock->socket_id = uasync_add_socket_t(instance->ua, e_sock->fd, socks_etcp_read_callback, NULL, NULL, "etcp_sock_socks", e_sock);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "socket %s tunneled via SOCKS5 UDP (proxy=%s:%u)", e_sock->name, gc->socks_host, gc->socks_port);
} else {
e_sock->socket_id = uasync_add_socket_t(instance->ua, e_sock->fd, etcp_connections_read_callback_socket, NULL, NULL, "etcp_sock", e_sock);
}
if (!e_sock->socket_id) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "Failed to register socket with uasync");
if (e_sock->socks_udp) { socks_udp_destroy(e_sock->socks_udp); e_sock->socks_udp = NULL; }
socket_close_wrapper(e_sock->fd);
queue_free(e_sock->links_queue);
instance->etcp_sockets = e_sock->next;
u_free(e_sock);
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Socket %p registered and active", e_sock);
return e_sock;
}
// Удаление UDP-сокета: снятие с uasync, закрытие fd, закрытие всех линков и освобождение.
void etcp_socket_remove(struct ETCP_SOCKET* conn) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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 && conn->instance) {
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");
}
if (conn->socks_udp) { socks_udp_destroy(conn->socks_udp); conn->socks_udp = NULL; }
if (conn->links_queue) {
struct ll_entry* entry;
while ((entry = conn->links_queue->head) != NULL) {
struct link_queue_entry* lqe = (struct link_queue_entry*)entry->data;
etcp_link_close(lqe->link);
}
queue_free(conn->links_queue);
conn->links_queue = NULL;
}
if (conn->instance) {
struct ETCP_SOCKET** pp = &conn->instance->etcp_sockets;
while (*pp && *pp != conn) pp = &(*pp)->next;
if (*pp) *pp = conn->next;
}
u_free(conn);
}
/* ── TCP socket management ── */
// TCP-сокет (запись ETCP_SOCKET с is_tcp=1): определение bind-адреса по конфигу.
// Сам listen/accept выполняет stcp_server — здесь только модель сокета.
struct ETCP_SOCKET* tcp_socket_add(struct UTUN_INSTANCE* instance, struct CFG_SERVER* server) {
if (!instance || !server) return NULL;
struct ETCP_SOCKET* ts = u_calloc(1, sizeof(struct ETCP_SOCKET));
if (!ts) { DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_socket_add: alloc failed"); return NULL; }
ts->instance = instance;
ts->is_tcp = 1;
ts->reality_enabled = server->reality_enabled;
ts->fd = SOCKET_INVALID;
ts->type = server->type;
if (server->name && server->name[0]) {
strncpy(ts->name, server->name, MAX_CONN_NAME_LEN - 1);
ts->name[MAX_CONN_NAME_LEN - 1] = '\0';
}
ts->sock_id = instance->next_socket_id++;
if (server->ip.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&server->ip;
memcpy(&ts->local_addr, sin, sizeof(*sin));
ts->interface_addr = ts->local_addr;
if (sin->sin_addr.s_addr == 0 && server->netif_index > 0) {
uint32_t if_ip = get_interface_ip_by_index(server->netif_index);
if (if_ip != 0) {
struct sockaddr_in* if_sin = (struct sockaddr_in*)&ts->interface_addr;
if_sin->sin_family = AF_INET;
if_sin->sin_addr.s_addr = if_ip;
if_sin->sin_port = sin->sin_port;
}
}
if (server->type == CFG_SERVER_TYPE_PUBLIC && sin->sin_addr.s_addr == 0) {
struct sockaddr_storage remote; memset(&remote, 0, sizeof(remote));
struct sockaddr_in* rs = (struct sockaddr_in*)&remote;
rs->sin_family = AF_INET; rs->sin_port = htons(53);
inet_pton(AF_INET, "8.8.8.8", &rs->sin_addr);
struct sockaddr_storage local;
if (get_outgoing_local_ip(&server->ip, &remote, &local) == 0) {
ts->interface_addr = local;
((struct sockaddr_in*)&ts->interface_addr)->sin_port = sin->sin_port;
}
}
} else if (server->ip.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&server->ip;
memcpy(&ts->local_addr, sin6, sizeof(*sin6));
ts->interface_addr = ts->local_addr;
if (memcmp(&sin6->sin6_addr, &in6addr_any, sizeof(struct in6_addr)) == 0 && server->netif_index > 0) {
uint8_t v6addr[16];
int got = get_interface_ipv6_addr_nl(server->netif_index, 1, v6addr);
if (got != 0) got = get_interface_ipv6_addr_nl(server->netif_index, 0, v6addr);
if (got == 0) {
struct sockaddr_in6* if6 = (struct sockaddr_in6*)&ts->interface_addr;
if6->sin6_family = AF_INET6; memcpy(&if6->sin6_addr, v6addr, 16);
if6->sin6_port = sin6->sin6_port;
}
}
if (server->type == CFG_SERVER_TYPE_PUBLIC && memcmp(&sin6->sin6_addr, &in6addr_any, sizeof(struct in6_addr)) == 0) {
uint16_t v6if = (server->netif_index > 0) ? server->netif_index : get_default_route_netif_index(AF_INET6);
if (v6if > 0) {
uint8_t v6addr[16];
int got = get_interface_ipv6_addr_nl(v6if, 1, v6addr);
if (got != 0) got = get_interface_ipv6_addr_nl(v6if, 0, v6addr);
if (got == 0) {
struct sockaddr_in6* if6 = (struct sockaddr_in6*)&ts->interface_addr;
if6->sin6_family = AF_INET6; memcpy(&if6->sin6_addr, v6addr, 16);
}
}
if (ts->interface_addr.ss_family == 0) {
struct sockaddr_storage remote; memset(&remote, 0, sizeof(remote));
struct sockaddr_in6* rs6 = (struct sockaddr_in6*)&remote;
rs6->sin6_family = AF_INET6; rs6->sin6_port = htons(53);
inet_pton(AF_INET6, "2001:4860:4860::8888", &rs6->sin6_addr);
if (get_outgoing_local_ip(&server->ip, &remote, &ts->interface_addr) == 0)
((struct sockaddr_in6*)&ts->interface_addr)->sin6_port = sin6->sin6_port;
else
ts->interface_addr = ts->local_addr;
}
}
}
ts->next = instance->etcp_sockets;
instance->etcp_sockets = ts;
{ char loc_str[64] = "none", if_str[64] = "none";
if (ts->local_addr.ss_family) snprintf(loc_str, sizeof(loc_str), "%s", sockaddr_storage_to_str(&ts->local_addr).str);
if (ts->interface_addr.ss_family) snprintf(if_str, sizeof(if_str), "%s", sockaddr_storage_to_str(&ts->interface_addr).str);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_add: %s type=%s sock_id=%u local=%s iface=%s",
ts->name, server_type_str(ts->type), ts->sock_id, loc_str, if_str); }
return ts;
}
// Удаление TCP-сокета из списка instance->etcp_sockets.
void tcp_socket_remove(struct ETCP_SOCKET* sock) {
if (!sock || !sock->instance) return;
struct UTUN_INSTANCE* inst = sock->instance;
struct ETCP_SOCKET** pp = &inst->etcp_sockets;
while (*pp && *pp != sock) pp = &(*pp)->next;
if (*pp) *pp = sock->next;
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "tcp_socket_remove: %s sock_id=%u", sock->name, sock->sock_id);
u_free(sock);
}
// Создание нового ETCP_LINK: аллокация, выделение local_link_id, инициализация BBR/keepalive,
// вставка в список коннекта и links_queue сокета. Для outbound UDP сразу запускает handshake.
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_ETCP, "");
if (!etcp) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "null etcp"); return NULL; }
if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] link_new on deleted conn", etcp->log_name); return NULL; }
int is_tcp = (conn == NULL) || conn->is_tcp;
struct ETCP_LINK* link = u_calloc(1, sizeof(struct ETCP_LINK));
if (!link) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "calloc failed - out of memory or pool exhausted");
return NULL;
}
link->conn = conn;
link->etcp = etcp;
link->is_server = is_server;
link->is_tcp = (uint8_t)is_tcp;
int mtu;
mtu = conn ? conn->mtu : STCP_MAX_MSG_SIZE;
if (is_tcp) { mtu = STCP_MAX_MSG_SIZE; link->mtu_local = mtu; link->mtu = mtu; link->inflight_lim_bytes = mtu * 4; link->handshake_minsize = 100; link->handshake_maxsize = mtu; }
else { if (mtu == 0) mtu = 1280; if (mtu > PACKET_DATA_MAX_MTU) mtu = PACKET_DATA_MAX_MTU; link->mtu_local = mtu; link->mtu = mtu; link->inflight_lim_bytes = mtu * 4; link->handshake_minsize = 100; link->handshake_maxsize = mtu; }
/* mtu init moved above */
etcp_update_mtu(etcp);
link->initialized = 0;
link->init_timer = NULL;
link->init_timeout = 0;
link->init_retry_count = 0;
link->link_status = 0;
link->send_hook = NULL;
link->send_hook_ctx = NULL;
// 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->ka_pinned = 0;
/* inflight_lim_bytes set above */
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_ETCP, "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;
if (remote_addr) memcpy(&link->remote_addr, remote_addr, sizeof(struct sockaddr_storage));
if (conn && conn->interface_addr.ss_family) memcpy(&link->local_bound_addr, &conn->interface_addr, sizeof(struct sockaddr_storage));
else memset(&link->local_bound_addr, 0, sizeof(link->local_bound_addr));
link->last_recv_local_time = get_time_tb(); // Initialize to prevent immediate timeout
link->total_retransmissions = 0;
if (!is_tcp && insert_link_queue(conn, link) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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_ETCP, "NEW link initialized on etcp=[%s] link=%p socket=%s id=%d is_server=%d mtu=%d is_tcp=%d",
etcp->log_name, link, conn ? conn->name : "tcp", link->local_link_id, link->is_server, link->mtu, link->is_tcp);
etcp_fire_link_status_cbk(link, -1, 0);
if (!is_tcp && is_server == 0) {
etcp_link_enter_init(link);
}
return link;
}
// Установка нового лимита inflight с клампингом в [INFLIGHT_LIM_MIN, max_inflight] и уведомлением коннекта.
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);
}
// Обнуляет last_link для всех inflight-записей ссылающихся на dead_link.
// Вызывается перед u_free(link) чтобы etcp_ack_recv не упал на use-after-free.
static void etcp_inflight_nullify_link(struct ll_queue* q, struct ETCP_LINK* dead_link) {
if (!q) return;
struct ll_entry* entry = q->head;
while (entry) {
struct INFLIGHT_PACKET* inf = (struct INFLIGHT_PACKET*)entry;
if (inf->last_link == dead_link) inf->last_link = NULL;
entry = entry->next;
}
}
// Закрытие линка: отмена таймеров (init/shaper/keepalive/tcp), удаление из списков и очереди,
// зануление ссылок inflight (защита от use-after-free) и освобождение памяти.
void etcp_link_close(struct ETCP_LINK* link) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!link) return;
if (link->etcp && link->etcp->last_rr_link == link) link->etcp->last_rr_link = NULL;
if (link->is_tcp) {
if (link->tcp_reconnect_timer) { uasync_cancel_timeout(link->etcp->instance->ua, link->tcp_reconnect_timer); link->tcp_reconnect_timer = NULL; }
if (link->tcp_link) { stcp_link_set_on_close(link->tcp_link, NULL, NULL); stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
}
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);
etcp_inflight_nullify_link(link->etcp->input_wait_ack, link);
etcp_inflight_nullify_link(link->etcp->input_send_q, link);
u_free(link->bbr);
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;
}
remove_link_from_queue(link);
etcp_conn_on_inflight_lim_changed(link->etcp);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d closed: addr=%s rcvd=%zub ack=%llub infl=%ub/%upkt",
link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str, link->total_decrypted,
(unsigned long long)link->acked_bytes, link->inflight_bytes, link->inflight_packets);
etcp_fire_link_status_cbk(link, link->link_state, link->link_status);
etcp_inflight_nullify_link(link->etcp->input_wait_ack, link);
etcp_inflight_nullify_link(link->etcp->input_send_q, link);
etcp_on_link_down(link->etcp, link);
u_free(link->bbr);
u_free(link);
}
// ====== TCP link helpers ======
// Поднять TCP-линк в ready: синхронизация reset-эпохи (gop/rid), keepalive, ed25519,
// mark initialized + link_status=1, уведомить коннект и loadbalancer.
void etcp_link_enter_ready_tcp(struct ETCP_LINK *link) {
if (!link || !link->etcp) return;
struct ETCP_CONN *etcp = link->etcp;
/* reinit только при смене reset_id пира — синхронизация через etcp_conn_apply_peer_reset_id */
if (link->tcp_link) {
uint64_t peer_rid = stcp_link_get_peer_reset_id(link->tcp_link);
etcp_conn_apply_peer_reset_id(etcp, peer_rid);
if (etcp->state == 2) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "enter_ready_tcp: conn DELETED after apply_reset, link=%p", (void*)link); return; }
link->peer_device_type = stcp_link_get_peer_device_type(link->tcp_link);
{ uint16_t peer_ka = stcp_link_get_peer_keepalive_interval(link->tcp_link);
link->keepalive_interval = negotiate_keepalive(etcp->instance->client_type,
etcp->instance->keepalive_interval, link->peer_device_type, peer_ka); }
}
int old_state = link->link_state;
link->initialized = 1; link->link_state = LINK_STATE_CONNECTED; link->link_status = 1;
link->recv_keepalive = 1;
link->last_recv_local_time = get_time_tb();
if (!link->mtu_remote) link->mtu_remote = link->mtu;
etcp->tcp_link_count++;
if (link->tcp_link) { const uint8_t* ed = stcp_link_get_peer_ed25519_pubkey(link->tcp_link); if (ed) memcpy(etcp->peer_ed25519_pubkey, ed, SC_PUBKEY_SIZE); }
if (etcp->initialized == 0) etcp_conn_ready(etcp);
if (etcp->state == 2) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "enter_ready_tcp: conn DELETED after conn_ready, link=%p", (void*)link); return; }
etcp_link_send_keepalive(link);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
etcp_fire_link_status_cbk(link, old_state, link->link_status);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d UP (mtu=%d init=%d up=%d tcp_links=%d)",
etcp->log_name, link->local_link_id, link->mtu,
etcp->initialized, etcp->links_up, etcp->tcp_link_count);
}
// Повторная попытка STCP-подключения с экспоненциальным backoff (до 30 с).
static void tcp_link_reconnect_cb(void *arg) {
struct ETCP_LINK *link = (struct ETCP_LINK *)arg;
if (!link || !link->etcp || !link->etcp->instance) return;
link->tcp_reconnect_timer = NULL;
if (link->tcp_reconnect_delay_ms == 0) link->tcp_reconnect_delay_ms = 1000;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d reconnect attempt (delay=%ums)", link->etcp->log_name, link->local_link_id, link->tcp_reconnect_delay_ms);
uint16_t port = ntohs(((struct sockaddr_in *)&link->remote_addr)->sin_port);
struct stcp_link *sl = stcp_link_connect(link, &link->remote_addr, port);
if (!sl) { link->tcp_reconnect_delay_ms *= 2; if (link->tcp_reconnect_delay_ms > 30000) link->tcp_reconnect_delay_ms = 30000;
link->tcp_reconnect_timer = uasync_set_timeout(link->etcp->instance->ua, (int)(link->tcp_reconnect_delay_ms * 10), link, tcp_link_reconnect_cb, "tcp_rct"); return; }
link->tcp_link = sl;
stcp_link_set_on_close(sl, tcp_link_close_cb, link);
}
// Исходящее TCP-подключение линка: stcp_link_connect + установка close-callback.
void etcp_tcp_link_start_connect(struct ETCP_LINK *link, struct sockaddr_storage *addr, uint16_t port) {
if (!link || !link->etcp || !addr) return;
/* link->conn опционален: NULL = без локального bind (ОС выбирает source IP/интерфейс).
* Задан — bind к interface_addr этого сокета (выбор интерфейса/ip). */
memcpy(&link->remote_addr, addr, sizeof(*addr));
struct stcp_link *sl = stcp_link_connect(link, &link->remote_addr, port);
if (!sl) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_tcp_link_start_connect: stcp_link_connect failed"); return; }
link->tcp_link = sl;
if (link->is_server == 0) stcp_link_set_on_close(sl, tcp_link_close_cb, link);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d → stcp_link_connect %s:%u rc=%p bind=%s",
link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(addr).str, port, (void*)sl,
link->conn ? sockaddr_storage_to_str(&link->conn->interface_addr).str : "auto");
}
// Запланировать реконнект TCP-линка после обрыва (через tcp_reconnect_timer).
void etcp_tcp_link_start_reconnect(struct ETCP_LINK *link) {
if (!link || !link->etcp || !link->etcp->instance || link->is_server) return;
if (link->etcp->state == 2) return;
if (link->tcp_reconnect_timer) return;
if (link->tcp_link) { stcp_link_close(link->tcp_link); link->tcp_link = NULL; }
link->link_state = LINK_STATE_HANDSHAKE;
if (link->etcp->tcp_link_count > 0) link->etcp->tcp_link_count--;
link->tcp_reconnect_timer = uasync_set_timeout(link->etcp->instance->ua, (int)(link->tcp_reconnect_delay_ms * 10), link, tcp_link_reconnect_cb, "tcp_rct");
}
// Callback закрытия STCP-линка: серверный линк закрываем, клиентский — отправляем на реконнект.
static void tcp_link_close_cb(struct stcp_link *sl, int err, void *arg) {
struct ETCP_LINK *link = (struct ETCP_LINK *)arg;
if (!link || !link->etcp) return;
if (link->etcp->state == 2) return;
/* Сохраняем до каскада: etcp_on_link_down может закрыть conn и освободить link */
int old_state = link->link_state;
int is_server = link->is_server;
int local_link_id = link->local_link_id;
struct ETCP_CONN *etcp = link->etcp;
stcp_link_close(sl);
link->tcp_link = NULL;
etcp_fire_link_status_cbk(link, old_state, 0);
if (is_server) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP server link %d down err=%d, closing", etcp->log_name, local_link_id, err);
etcp_on_link_down(etcp, link);
if (etcp->state != 2) {
etcp_link_close(link);
} else {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP server link %d already closed by conn teardown, skip",
etcp->log_name, local_link_id);
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP link %d down err=%d, scheduling reconnect", etcp->log_name, local_link_id, err);
etcp_on_link_down(etcp, link);
if (etcp->state != 2) etcp_tcp_link_start_reconnect(link);
}
}
// Отправка ETCP_DGRAM по TCP: сборка буфера (timestamp+flag_up+data) и передача в stcp_link_send.
static int etcp_tcp_send(struct ETCP_DGRAM* dgram) {
struct ETCP_LINK *link = dgram->link;
if (!link->tcp_link || !stcp_link_is_ready(link->tcp_link)) return -1;
link->pkt_sent_since_keepalive = 1;
dgram->flag_up = 1;
dgram->timestamp = get_current_timestamp();
size_t total_len = 3 + dgram->data_len;
uint8_t *buf = u_malloc(total_len);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] TCP send: malloc %zu failed", link->etcp->log_name, total_len); return -1; }
memcpy(buf, &dgram->timestamp, 2);
buf[2] = dgram->flag_up;
if (dgram->data_len) memcpy(buf + 3, dgram->data, dgram->data_len);
int rc;
if (link->send_hook) rc = (int)link->send_hook(0, buf, total_len, NULL, 0, link, link->send_hook_ctx);
else if (link->tcp_link) rc = stcp_link_send(link->tcp_link, buf, total_len) == 0 ? (int)dgram->data_len : -1;
else rc = -1;
u_free(buf);
if (rc > 0) link->total_encrypted += (size_t)rc;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] TCP send: link=%d total=%zu dlen=%d rc=%d via=%s",
link->etcp->log_name, link->local_link_id, total_len, dgram->data_len, rc,
link->send_hook ? "hook" : (link->tcp_link ? "stcp" : "NONE"));
return rc;
}
// Единая точка отправки UDP: SOCKS5-релей (если сокет proxied) → send_hook (dummynet) → socket_sendto.
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->conn && link->conn->socks_udp)
return socks_udp_sendto(link->conn->socks_udp, buf, len, (const struct sockaddr_storage*)addr);
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);
}
// Зашифровать и отправить пакет: TCP → etcp_tcp_send, UDP → AES-CCM + socket_sendto.
int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
if (!dgram || !dgram->link) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Null pointer"); return -1; }
if (dgram->link->is_tcp) return etcp_tcp_send(dgram);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%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;
int udp_ov = (dgram->link->remote_addr.ss_family == AF_INET6) ? IP_UDP_OVERHEAD_V6 : IP_UDP_OVERHEAD_V4;
int wire_ov = udp_ov + SC_ENCRYPT_OVERHEAD;
if (len<0 || len + (int)dgram->noencrypt_len > (int)(dgram->link->mtu - wire_ov)) { dgram->link->send_errors++; errcode=1; goto es_err; }
uint8_t enc_buf[PACKET_DATA_MAX_MTU];// макс запись = data_len + 32 <= mtu - wire_ov + 32 = 1468
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_ETCP, "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_ov)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "packet too long enc=%zu ne=%d mtu=%d", enc_buf_len, dgram->noencrypt_len, dgram->link->mtu);
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);
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_ETCP, "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_ETCP, "error %d", errcode);
return -1;
}
// Шифрует и отправляет готовый PING/PONG-датаграмм без линка (сырой sendto по сокету).
// Для proxied-сокета идёт через SOCKS5 UDP-релей (dst = addr пира).
static int etcp_send_ping_raw(struct ETCP_DGRAM* dgram, struct ETCP_SOCKET* e_sock, sc_context_t* sc, const struct sockaddr_storage* addr) {
if (!dgram || !e_sock || !sc || !addr) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Null pointer in ping send");
return -1;
}
int len = dgram->data_len - dgram->noencrypt_len;
uint8_t enc_buf[1600];
if (len < 0 || len + (int)dgram->noencrypt_len + SC_ENCRYPT_OVERHEAD > (int)sizeof(enc_buf)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ping packet too large len=%d ne=%d", len, dgram->noencrypt_len);
return -1;
}
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_ETCP, "encryption failed for ping");
return -1;
}
if (enc_buf_len + dgram->noencrypt_len > (size_t)(PACKET_DATA_SIZE - IP_UDP_OVERHEAD_V4)) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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);
size_t total = enc_buf_len + dgram->noencrypt_len;
ssize_t sent;
if (e_sock->socks_udp) {
sent = socks_udp_sendto(e_sock->socks_udp, enc_buf, total, addr);
if (sent < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ping via socks failed addr=%s len=%zu", sockaddr_storage_to_str(addr).str, total);
return -1;
}
return (int)sent;
}
socklen_t addr_len = (addr->ss_family == AF_INET) ? sizeof(struct sockaddr_in) : sizeof(struct sockaddr_in6);
sent = socket_sendto(e_sock->fd, enc_buf, total, (struct sockaddr*)addr, addr_len);
if (sent < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "sendto failed for ping, err=%d addr=%s fd=%d len=%zu", socket_get_error(), sockaddr_storage_to_str(addr).str, e_sock->fd, total);
return -1;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "ping sendto succeeded to %s sent=%zd bytes fd=%d", sockaddr_storage_to_str(addr).str, sent, e_sock->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;
}
// Таймаут пинга: снять из pending_pings и вызвать callback с success=0.
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);
}
// Отправка UDP-пинга через конкретный сокет: сборка PING-кодограммы (nonce, user_data,
// обфусцированный pubkey), регистрация в pending_pings с таймаутом.
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,
uint8_t flags) {
if (!instance || !e_sock || !peer_pubkey_bin || !addr || timeout_ms <= 0 || !cb) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "bad args");
return -1;
}
if (e_sock->is_tcp || e_sock->fd == SOCKET_INVALID) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ping on non-UDP socket name=%s is_tcp=%d fd=%d",
e_sock->name, e_sock->is_tcp, (int)e_sock->fd);
return -2;
}
if (user_data_len > PACKET_DATA_SIZE - 24 - 2 - SC_PUBKEY_ENC_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "user_data too long");
return -2;
}
struct PING_CONTEXT* ctx = u_malloc(sizeof(struct PING_CONTEXT));
if (!ctx) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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_ETCP, "malloc user_data");
return -3;
}
memcpy(ctx->user_data, user_data, user_data_len);
ctx->user_data_len = user_data_len;
}
memcpy(ctx->peer_pubkey, peer_pubkey_bin, SC_PUBKEY_SIZE);
struct ETCP_DGRAM* dgram = u_malloc(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
if (!dgram) {
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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;
*p++ = flags; // flags
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_ETCP, "ping send nonce=%016llx timeout=%d ulen=%zu",
(unsigned long long)ctx->nonce, timeout_ms, user_data_len);
int send_rc = etcp_send_ping_raw(dgram, e_sock, &sc, addr);
u_free(dgram);
if (send_rc < 0) {
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
return -6;
}
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;
}
// Отправка UDP-пинга: подбор сокета по family адреса, делегирование в etcp_send_ping_to_socket.
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_ETCP, "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->is_tcp || e_sock->local_addr.ss_family != addr->ss_family)) e_sock = e_sock->next;
if (!e_sock) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no UDP socket for addr_family=%d", addr->ss_family);
return -2;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "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, 0);
}
struct tcp_ping_adapter {
etcp_ping_callback_t cb;
void* arg;
};
// Адаптер: приводит callback STCP-пинга к типу etcp_ping_callback_t.
static void tcp_ping_cb_adapter(int success, uint16_t rtt, void* arg) {
struct tcp_ping_adapter* a = (struct tcp_ping_adapter*)arg;
a->cb(success, rtt, a->arg, 0, NULL, 0);
u_free(a);
}
// TCP-пинг: STCP-хендшейк с пиром и замер RTT по handshake (через stcp_ping_send).
int etcp_send_tcp_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bin,
const struct sockaddr_storage* addr, int timeout_ms,
etcp_ping_callback_t cb, void* user_arg) {
if (!instance || !peer_pubkey_bin || !addr || timeout_ms <= 0 || !cb) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "bad args: inst=%p pubkey=%p addr=%p timeout=%d cb=%p",
(void*)instance, (void*)peer_pubkey_bin, (void*)addr, timeout_ms, (void*)(uintptr_t)cb);
return -1;
}
char addr_str[INET6_ADDRSTRLEN];
uint16_t port;
if (addr->ss_family == AF_INET) {
const struct sockaddr_in* sin = (const struct sockaddr_in*)addr;
snprintf(addr_str, sizeof(addr_str), "%s", ip_to_str(&sin->sin_addr, AF_INET).str);
port = ntohs(sin->sin_port);
} else if (addr->ss_family == AF_INET6) {
const struct sockaddr_in6* sin6 = (const struct sockaddr_in6*)addr;
snprintf(addr_str, sizeof(addr_str), "%s", ip_to_str(&sin6->sin6_addr, AF_INET6).str);
port = ntohs(sin6->sin6_port);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "unsupported addr family=%d", addr->ss_family);
return -2;
}
struct tcp_ping_adapter* a = u_malloc(sizeof(struct tcp_ping_adapter));
if (!a) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "malloc adapter"); return -3; }
a->cb = cb; a->arg = user_arg;
struct socks_cfg socks_buf;
const struct socks_cfg* socks = NULL;
if (instance->config && instance->config->global.socks_enabled &&
instance->config->global.socks_host[0] && instance->config->global.socks_port) {
memset(&socks_buf, 0, sizeof(socks_buf));
strncpy(socks_buf.host, instance->config->global.socks_host, sizeof(socks_buf.host) - 1);
socks_buf.port = instance->config->global.socks_port;
strncpy(socks_buf.user, instance->config->global.socks_username, sizeof(socks_buf.user) - 1);
strncpy(socks_buf.pass, instance->config->global.socks_password, sizeof(socks_buf.pass) - 1);
socks = &socks_buf;
}
struct stcp_client* cli = stcp_ping_send(instance->ua, addr_str, port, &instance->my_keys, peer_pubkey_bin,
instance->my_ed25519_pubkey, instance->client_type,
instance->keepalive_interval, timeout_ms, tcp_ping_cb_adapter, a, socks);
if (!cli) { u_free(a); return -4; }
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "tcp ping to %s:%u timeout=%d", addr_str, (unsigned)port, timeout_ms);
return 0;
}
// === Helpers extracted from etcp_connections_read_callback_socket ===
// Обработка входящего PING: обновление RTT (если флаг SEND_RTT), ответ PONG с эхо-данными.
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 < 23) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "PING too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str);
return 7;
}
uint8_t flags = pkt->data[1];
uint64_t peer_id = be64toh(*(uint64_t*)(pkt->data + 2));
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 10));
uint16_t ulen = be16toh(*(uint16_t*)(pkt->data + 18));
if (ulen > pkt->data_len - 20) ulen = (uint16_t)(pkt->data_len - 20);
const uint8_t* udata = (ulen > 0) ? (pkt->data + 20) : NULL;
if ((flags & ETCP_PING_FLAG_SEND_RTT) && ulen >= 2) {
uint16_t rtt_val = be16toh(*(uint16_t*)udata);
topo_node_ping_update_rtt(e_sock->instance->topo_groups, peer_id, rtt_val);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "PING rtt=%u from peer=%016llx", (unsigned)rtt_val, (unsigned long long)peer_id);
}
uint8_t pong_flags = 0;
if (e_sock->instance->topo_groups && topo_node_ping_request_cbk(e_sock->instance->topo_groups, peer_id))
pong_flags |= ETCP_PING_FLAG_WANT_RTT;
struct ETCP_DGRAM* resp = u_malloc(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE);
if (resp) {
resp->link = NULL;
resp->noencrypt_len = SC_PUBKEY_ENC_SIZE;
uint8_t* p = resp->data;
*p++ = ETCP_PONG;
*p++ = pong_flags;
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_ETCP, "PONG send nonce=%016llx flags=%02x to=%s fd=%d",
(unsigned long long)nonce, (unsigned)pong_flags, sockaddr_storage_to_str(addr).str, e_sock->fd);
etcp_send_ping_raw(resp, e_sock, &resp_sc, addr);
}
u_free(resp);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
// Пустой callback для RTT-ответа (получателю RTT-пинга ответ не важен).
static void rtt_send_ping_cb(int success, uint16_t rtt, void* arg, uint64_t nonce,
const uint8_t* resp_data, size_t resp_data_len) {
(void)success; (void)rtt; (void)arg; (void)nonce; (void)resp_data; (void)resp_data_len;
}
// Обработка входящего PONG: поиск pending_pings по nonce, вызов callback с RTT и данными.
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 < 23) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "PONG too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str);
return 7;
}
uint8_t flags = pkt->data[1];
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 10));
uint16_t ulen = 0;
const uint8_t* udata = NULL;
if (pkt->data_len >= 20) {
ulen = be16toh(*(uint16_t*)(pkt->data + 18));
if (ulen > 0 && pkt->data_len >= 20 + ulen) {
udata = pkt->data + 20;
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "PONG recv nonce=%016llx flags=%02x data_len=%u from=%s socket=%s",
(unsigned long long)nonce, (unsigned)flags, (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_ETCP, "PONG matched nonce=%016llx rtt=%u want_rtt=%d",
(unsigned long long)nonce, (unsigned)rtt, (flags & ETCP_PING_FLAG_WANT_RTT) ? 1 : 0);
uint64_t pong_peer_id = be64toh(*(uint64_t*)(pkt->data + 2));
topo_node_ping_update_rtt(e_sock->instance->topo_groups, pong_peer_id, rtt);
ctx->cb(1, rtt, ctx->arg, nonce, udata, ulen);
if ((flags & ETCP_PING_FLAG_WANT_RTT) && ctx->peer_pubkey[0] != 0) {
uint8_t rtt_buf[2];
rtt_buf[0] = (uint8_t)(rtt >> 8);
rtt_buf[1] = (uint8_t)(rtt & 0xFF);
etcp_send_ping_to_socket(e_sock->instance, e_sock, ctx->peer_pubkey, addr, 1000,
rtt_send_ping_cb, NULL, rtt_buf, 2, ETCP_PING_FLAG_SEND_RTT);
}
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
break;
}
prev = ctx;
ctx = ctx->next;
}
if (!found) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "PONG nonce=%016llx NOT FOUND in pending (timeout?)",
(unsigned long long)nonce);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
// Сервер: собрать и отправить INIT_RESPONSE (mtu, link_id, NAT-адрес клиента, ed25519),
// провести NAT-детекцию, поднять линк и запустить keepalive.
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_ETCP, "send init_response with reset");
resp->code = ETCP_INIT_RESPONSE; // 0x03 - with reset
} else {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "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);
*(uint64_t*)resp->reset_id = htobe64(conn->reset_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->nat_det) {
nat_detection_send_nat_info(link->etcp->instance->nat_det, link->etcp,
link->remote_socket_id,
link->nat_ip, link->nat_port, NAT_TYPE_DIRECT);
}
DEBUG_INFO(DEBUG_CATEGORY_NAT, "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=%s nat=%s if_addr=%s local_addr=%s",
e_sock->name, server_type_str(e_sock->type), nat_type_str(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) {
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, 0);
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);
resp->device_type = e_sock->instance->client_type;
*(uint16_t*)resp->keepalive = htobe16(e_sock->instance->keepalive_interval);
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 wire_ov = WIRE_OVERHEAD(link->remote_addr.ss_family);
int max_data = link->mtu - wire_ov;
int pad_range = (int)link->handshake_maxsize - (int)link->handshake_minsize;
if (pad_range <= 0) pad_range = 1;
int s = rand() % pad_range + link->handshake_minsize;
if (s > (int)(link->mtu)) s = (int)(link->mtu);
if (s < 0) s = 0;
int to_add = s - xoffset - wire_ov;
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; }
if (xoffset + to_add > max_data) to_add = max_data - 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_INFO(DEBUG_CATEGORY_ETCP, "[%s] INIT_RESPONSE sent (link=%d, mtu=%d, dst=%s)",
link->etcp->log_name, link->local_link_id, link->mtu_local,
sockaddr_storage_to_str(&link->remote_addr).str);
etcp_encrypt_send(pkt);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
link->initialized = 1;
{ int old_state = link->link_state; link->link_state = LINK_STATE_CONNECTED; etcp_fire_link_status_cbk(link, old_state, link->link_status); }
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_ETCP, "[%s] Connection established (socket=%s, link=%d, addr=%s)", link->etcp->log_name, e_sock->name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str);
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Link %d UP (server, mtu=%d, addr=%s)", link->etcp->log_name, link->local_link_id, link->mtu_local, sockaddr_storage_to_str(&link->remote_addr).str);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
// Restart NAT check after link is up (e.g. after address change or reinit)
if (link->etcp->instance->nat_det) nat_detection_link_ready(link->etcp->instance->nat_det, link);
}
// Клиент: обработка INIT_RESPONSE — MTU, link_id, NAT-адрес, ed25519, keepalive, подъём линка.
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 (!e_sock) return -1;
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);
uint64_t resp_reset_id = be64toh(*(uint64_t*)resp->reset_id);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE rid=%016llx", (unsigned long long)resp_reset_id);
etcp_conn_apply_peer_reset_id(link->etcp, resp_reset_id);
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;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] INIT_RESPONSE received (link=%d, mtu=%d, rid=%016llx)", link->etcp->log_name, link->remote_link_id, link->mtu, (unsigned long long)resp_reset_id);
// 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=%s new_ip=%s new_port=%u",
e_sock->name, nat_type_str(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=%s",
e_sock->name, nat_type_str(e_sock->nat_type));
}
}
} else {
// Legacy format without NAT info
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Received legacy INIT_RESPONSE without NAT info",
link->etcp->log_name);
}
memcpy(link->etcp->peer_ed25519_pubkey, resp->ed25519_pubkey, SC_PUBKEY_SIZE);
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "[%s] Received Ed25519 pubkey from peer", link->etcp->log_name);
link->peer_device_type = resp->device_type;
{ uint16_t server_ka = (resp->keepalive[0]<<8) | resp->keepalive[1];
link->keepalive_interval = negotiate_keepalive(link->etcp->instance->client_type,
link->etcp->instance->keepalive_interval, resp->device_type, server_ka); }
link->etcp->peer_node_id = server_node_id;
etcp_update_log_name(link->etcp);
link->initialized = 1;// получен init response (client)
{ int old_state = link->link_state; link->link_state = LINK_STATE_CONNECTED; etcp_fire_link_status_cbk(link, old_state, link->link_status); }
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);
}
loadbalancer_link_ready(link);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Link %d UP (client, mtu=%d, addr=%s)", link->etcp->log_name, link->local_link_id, link->mtu, sockaddr_storage_to_str(&link->remote_addr).str);
// 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;
}
// Главный приёмник UDP-кодограмм: normal/init decrypt, диспетчеризация PING/PONG/INIT,
// создание/переиспользование линков и коннектов, обработка коллизий.
// Приём прямого UDP-пакета: recvfrom → etcp_process_packet.
void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
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;
}
etcp_process_packet(e_sock, data, recv_len, addr);
}
// Обработка принятого сырого UDP-пакета: расшифровка и диспетчеризация.
// Выделена из read_callback, чтобы один путь обслуживал прямой приём и приём через
// SOCKS5 UDP ASSOCIATE (socks_etcp_read_callback подставляет реальный src пира).
// !!!!!! DANGER: в этой функции ПРЕДЕЛЬНАЯ АККУРАТНОСТЬ !!!!!
// НЕ РУИНИТЬ (uint8_t*)&pkt->timestamp - это правильно !!!!
//
// Ошибки функции (errorcode):
// 1 - пакет слишком маленький для init (< SC_PUBKEY_SIZE)
// 2 - не удалось установить peer public key при init
// 3 - не удалось расшифровать init пакет
// 4 - не init/ping пакет (неверный код)
// 5 - коллизия peer ID и ключей
// 6 - не удалось расшифровать обычный пакет
// 7 - слишком короткий пакет
// 8 - ключ не в списке allowed_keys
// 13 - переполнение при парсинге пакета
// 46 - расшифрованный пакет слишком маленький (< 3 байта)
// 55 - не удалось создать подключение
// 66 - не удалось создать линк
static void etcp_process_packet(struct ETCP_SOCKET* e_sock, uint8_t* data,
ssize_t recv_len, struct sockaddr_storage addr) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
// 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) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pkt_pool exhausted, dropping packet from %s", sockaddr_storage_to_str(&addr).str); return; }
size_t pkt_len=0;
int errorcode=0;
struct ETCP_LINK* link=etcp_link_find_by_addr(e_sock, &addr, 0);
// 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) {
DEBUG_TRACE(DEBUG_CATEGORY_CRYPTO, "NORM_DECRYPT_TRY: seskey=%02x%02x%02x%02x recv_len=%zd rx=%llu",
link->etcp->crypto_ctx.session_key[0], link->etcp->crypto_ctx.session_key[1],
link->etcp->crypto_ctx.session_key[2], link->etcp->crypto_ctx.session_key[3],
recv_len, (unsigned long long)link->etcp->crypto_ctx.rx_counter);
sc_status_t dec_rc = sc_decrypt(&link->etcp->crypto_ctx, data, recv_len, (uint8_t*)&pkt->timestamp, &pkt_len);
if (!dec_rc) {
int ec = etcp_packet_decrypted(e_sock, pkt, link, pkt_len);
if (ec) { errorcode = ec; goto ec_fr; }
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp: DECRYPT FAIL on existing link. rc=%d link=%p log=%s from=%s sess=%d link_state=%d keepalive=%d enc_errs=%zu my_pub=%016llx peer_pub=%016llx seskey=%02x%02x%02x%02x",
dec_rc, link, link->etcp->log_name, sockaddr_storage_to_str(&addr).str, 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,
link->etcp->crypto_ctx.session_key[0], link->etcp->crypto_ctx.session_key[1],
link->etcp->crypto_ctx.session_key[2], link->etcp->crypto_ctx.session_key[3]);
} else {
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "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) {
if (e_sock->pkt_format_errors < 2)
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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_DEBUG(DEBUG_CATEGORY_CRYPTO, "X25519 decrypt OK from %s", sockaddr_storage_to_str(&addr).str);
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "INIT_DECRYPT_TRY: seskey=%02x%02x%02x%02x recv_len=%zd strip40=%zd salt=%02x%02x%02x%02x%02x%02x%02x%02x enc_pub=%02x%02x%02x%02x",
sc.session_key[0], sc.session_key[1], sc.session_key[2], sc.session_key[3],
recv_len, recv_len - SC_PUBKEY_ENC_SIZE,
salt[0], salt[1], salt[2], salt[3], salt[4], salt[5], salt[6], salt[7],
encrypted_pubkey[0], encrypted_pubkey[1], encrypted_pubkey[2], encrypted_pubkey[3]);
if (sc_decrypt(&sc, data, recv_len - SC_PUBKEY_ENC_SIZE, (uint8_t*)&pkt->timestamp, &pkt_len)) {
#ifdef UTUN_HAVE_STANDBY
{ char who[64]; snprintf(who, sizeof(who), "addr=%s", sockaddr_storage_to_str(&addr).str); standby_log_rx("udp undecryptable", who); }
#endif
if (link) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO,
"packet undecryptable (normal+init fail) from=%s — existing link log=%s state=%d sess=%d",
sockaddr_storage_to_str(&addr).str,
link->etcp ? link->etcp->log_name : "null",
link->link_state,
link->etcp ? link->etcp->crypto_ctx.session_ready : -1);
} else if (e_sock->pkt_format_errors < 2) {
DEBUG_WARN(DEBUG_CATEGORY_CRYPTO,
"packet undecryptable (normal+init fail) from=%s — no link for this address",
sockaddr_storage_to_str(&addr).str);
}
errorcode=3;
goto ec_fr;
}
// INIT decryption succeeded - process packet
if (pkt_len<3) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "too short packet, from %s", sockaddr_storage_to_str(&addr).str);
errorcode=7;
goto ec_fr;
}
if (pkt_len<15) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "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 + 2));
if (code == ETCP_PING) {
#ifdef UTUN_HAVE_STANDBY
{ char who[64]; snprintf(who, sizeof(who), "addr=%s", sockaddr_storage_to_str(&addr).str); standby_log_rx("udp ping", who); }
#endif
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "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) {
#ifdef UTUN_HAVE_STANDBY
{ char who[64]; snprintf(who, sizeof(who), "addr=%s", sockaddr_storage_to_str(&addr).str); standby_log_rx("udp pong", who); }
#endif
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "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_ETCP, "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_ETCP, "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_ETCP, "INIT REQUEST too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(&addr).str);
errorcode=7;
goto ec_fr;
}
if (link)
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "INIT packet on existing link from=%s log=%s link_state=%d",
sockaddr_storage_to_str(&addr).str,
link->etcp ? link->etcp->log_name : "null", link->link_state);
// 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_ETCP, "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_ETCP, "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;
peer_id = be64toh(*(uint64_t*)req->node_id);
uint64_t reset_id = be64toh(*(uint64_t*)req->reset_id);
uint16_t mtu = be16toh(*(uint16_t*)req->mtu);
uint16_t src_port = be16toh(*(uint16_t*)req->src_port);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "INIT received: peer=0x%016llx mtu=%u link=%u sock=%u rid=%016llx src=%s",
(unsigned long long)peer_id, mtu, req->link_id, req->socket_id,
(unsigned long long)reset_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, 0); if (ol && ol->etcp && ol->etcp->conn_queue_entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)ol->etcp->conn_queue_entry->data;
if (ce->peer_node_id == 0) {
conn = ol->etcp;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] reusing unindexed outbound conn (q_key=0, peer=0x%016llx) for incoming INIT from %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_ETCP, "failed to create connection"); goto ec_fr; }
memcpy(&conn->crypto_ctx, &sc, sizeof(sc));
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "CRYPTO_CTX_INIT: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x my_priv=%02x%02x%02x%02x",
conn->log_name,
conn->crypto_ctx.session_key[0], conn->crypto_ctx.session_key[1],
conn->crypto_ctx.session_key[2], conn->crypto_ctx.session_key[3],
conn->crypto_ctx.peer_public_key[0], conn->crypto_ctx.peer_public_key[1],
conn->crypto_ctx.peer_public_key[2], conn->crypto_ctx.peer_public_key[3],
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[0] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[1] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[2] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[3] : 0);
conn->peer_node_id = peer_id;
etcp_update_log_name(conn);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "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);
}
else {// check keys если существующее подключение
if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) {
errorcode=5;
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node_changed: different pubkey from %s claiming peer=0x%016llx — conn=%s old_pub=%016llx new_pub=%016llx",
sockaddr_storage_to_str(&addr).str, (unsigned long long)peer_id, conn->log_name,
*(uint64_t*)conn->crypto_ctx.peer_public_key, *(uint64_t*)sc.peer_public_key);
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "peer key mismatch for node %016llx, firing node_changed callbacks", (unsigned long long)peer_id);
conn->callbacks_running = 1;
etcp_cbk_fire(conn, ETCP_CBK_EVENT_NODE_CHANGED);
conn->callbacks_running = 0;
goto ec_fr;
}// коллизия - peer id совпал а ключи разные.
}
if (conn->fin_wait) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] received INIT during fin_wait, clearing", conn->log_name);
conn->fin_wait = 0;
if (conn->fin_wait_clear_cb) { conn->fin_wait_clear_cb(conn, conn->fin_wait_clear_arg); conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL; }
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "INIT conn=%s new_conn=%d peer=0x%016llx state=%d links_up=%d links=%p", conn->log_name, new_conn, (unsigned long long)peer_id, conn->state, conn->links_up, (void*)conn->links);
// Check if link already exists (for CHANNEL_INIT recovery)
struct ETCP_LINK* existing_link = etcp_link_find_by_remote_id(conn, req->link_id);
/* Правила 1+2: найден по id, линк down, адрес сменился */
if (existing_link && existing_link->link_status == 0
&& !sockaddr_equal(&existing_link->remote_addr, &addr)) {
struct ETCP_LINK* L_addr = etcp_link_find_by_addr(e_sock, &addr, 0);
if (L_addr && L_addr != existing_link && L_addr->etcp == conn && L_addr->link_status == 0) {
/* Правило 2: своп remote_link_id, дальше работаем с L_addr */
uint8_t id_old = existing_link->remote_link_id;
uint8_t id_swap = L_addr->remote_link_id;
existing_link->remote_link_id = id_swap;
L_addr->remote_link_id = id_old;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] swap remote_link_id: id=%u @ %s <-> id=%u @ %s",
conn->log_name, id_old, sockaddr_storage_to_str(&existing_link->remote_addr).str,
id_swap, sockaddr_storage_to_str(&L_addr->remote_addr).str);
existing_link = L_addr;
} else if (!L_addr) {
/* Правило 1: перепривязать адрес */
etcp_link_update_remote_addr(existing_link, e_sock, &addr);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] remote addr updated: id=%u -> %s",
conn->log_name, existing_link->remote_link_id, sockaddr_storage_to_str(&addr).str);
} else {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] addr %s busy (up/foreign) — link id=%u left as-is",
conn->log_name, sockaddr_storage_to_str(&addr).str, existing_link->remote_link_id);
}
}
if (!existing_link) {
existing_link = etcp_link_find_by_addr(e_sock, &addr, 0);
if (existing_link && existing_link->etcp == conn) {
if (memcmp(conn->crypto_ctx.peer_public_key, sc.peer_public_key, SC_PUBKEY_SIZE)) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] link address match but pubkey mismatch, firing node_changed", conn->log_name);
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "pubkey mismatch on reused link for node %016llx", (unsigned long long)peer_id);
conn->callbacks_running = 1;
etcp_cbk_fire(conn, ETCP_CBK_EVENT_NODE_CHANGED);
conn->callbacks_running = 0;
errorcode = 67;
goto ec_fr;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] found existing outbound link by addr for incoming INIT, reusing link=%p",
conn->log_name, existing_link);
} else if (existing_link) {
struct ETCP_CONN* old_conn = existing_link->etcp;
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] conflicting link at %s belongs to %s, firing node_changed",
conn->log_name, sockaddr_storage_to_str(&addr).str, old_conn->log_name);
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "link address conflict for node %016llx, firing node_changed callbacks on %s",
(unsigned long long)peer_id, old_conn->log_name);
old_conn->callbacks_running = 1;
etcp_cbk_fire(old_conn, ETCP_CBK_EVENT_NODE_CHANGED);
old_conn->callbacks_running = 0;
errorcode = 67;
goto ec_fr;
} else {
existing_link = NULL;
}
}
uint8_t send_reset = 0;
if (existing_link && existing_link->etcp == conn) {// существующий линк
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%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_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] IP:port changed for remote_link_id=%d socket:[%s] — closing old link, creating new",
conn->log_name, req->link_id, e_sock->name);
etcp_link_close(link);
goto create_new_link;
}
// 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_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave",
conn->log_name, (unsigned long long)peer_id);
etcp_conn_reinit(conn, "collision yield");
send_reset = 1;
} else {
// Check if WE have an outbound (master) link on this conn
struct ETCP_LINK* ml = etcp_select_collision_link(conn, &addr);
if (ml) {
if (conn->instance->node_id < peer_id) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%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_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] COLLISION: peer smaller, yielding master, processing as slave",
conn->log_name);
if (ml->init_timer) { uasync_cancel_timeout(conn->instance->ua, ml->init_timer); ml->init_timer = NULL; }
}
if (req->code == ETCP_INIT_REQUEST) {
if (conn->got_initial_pkt) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] INIT_REQUEST on existing link: code=0x%02x got_init=%d → reinit",
conn->log_name, req->code, conn->got_initial_pkt);
etcp_conn_reinit(conn, "duplicate INIT");
}
send_reset = 0;
} else {
if (conn->got_initial_pkt == 0) send_reset = 1;
}
}
// Cancel existing timers
if (link->init_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
} else {
create_new_link:
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%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_ETCP, "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_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] INIT collision=1 from peer=%016llx — remote is master, becoming slave",
conn->log_name, (unsigned long long)peer_id);
etcp_conn_reinit(conn, "collision yield");
send_reset = 1;
} else {
// Check if WE have an outbound (master) link
struct ETCP_LINK* ml = etcp_select_collision_link(conn, &addr);
if (ml && conn->instance->node_id < peer_id) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%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;
}
if (conn->got_initial_pkt == 0) send_reset = 1;
}
uint16_t client_ka = (req->keepalive[0]<<8) | req->keepalive[1];
link->peer_device_type = req->device_type;
link->keepalive_interval = negotiate_keepalive(e_sock->instance->client_type, e_sock->instance->keepalive_interval, req->device_type, client_ka);
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);
// Применяем эпоху ресета пира ПОСЛЕ всей reinit-логики
etcp_conn_apply_peer_reset_id(conn, reset_id);
send_init_response(e_sock, pkt, link, conn, req, &addr, send_reset, req_src_ip, req_src_port, pkt_len);
return;
ec_fr:
e_sock->pkt_format_errors++;
if (e_sock->pkt_format_errors < 3 || (e_sock->pkt_format_errors % 500 == 0))
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "error %d, from %s (count=%zu)", 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;
}
// Обработка пакета, расшифрованного на установленном линке: keepalive-статус, INIT_RESPONSE,
// обычные данные → etcp_conn_input.
int etcp_packet_decrypted(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
struct ETCP_LINK* link, size_t pkt_len) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "Decrypt ok - normal pkt");
if (pkt_len<3) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "decrypted packet too small, size=%zu", pkt_len); return 46; }
pkt->data_len=pkt_len-3;
pkt->noencrypt_len=0;
pkt->link=link;
if (pkt->data_len && pkt->data[0] != ETCP_KEEPALIVE) {
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "[%s] decrypt: code=%02x dlen=%u plen=%zu recv=%d sock=%p",
link->etcp->log_name, pkt->data[0], pkt->data_len, pkt_len,
link->recv_keepalive, (void*)e_sock);
}
link->remote_keepalive = pkt->flag_up;
if (link->recv_keepalive != 1) {
link->recv_keepalive = 1;
DEBUG_DEBUG(DEBUG_CATEGORY_KEEPALIVE, "[%s] Link %d recv_keepalive restored to 1 (was 0) link_status=%d remote_ka=%d state=%d init=%d links_up=%d",
link->etcp->log_name, link->local_link_id, link->link_status, link->remote_keepalive, link->link_state, link->initialized, link->etcp->links_up);
}
int was_up = link->link_status;
link->link_status = link->is_tcp ? link->recv_keepalive : (link->remote_keepalive && link->recv_keepalive);
if (link->link_status != was_up)
etcp_fire_link_status_cbk(link, link->link_state, was_up);
if (link->link_status && !was_up && link->initialized) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d status popped UP: recv_ka=%d remote_ka=%d state=%d init=%d", link->etcp->log_name, link->local_link_id, link->recv_keepalive, link->remote_keepalive, link->link_state, 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;
link->total_decrypted += pkt->data_len;
uint8_t pkt_code = pkt->data[0];
#ifdef UTUN_HAVE_STANDBY
{
char who[256];
snprintf(who, sizeof(who), "%s addr=%s", link->etcp->log_name,
sockaddr_storage_to_str(&link->remote_addr).str);
int is_ka = (pkt_code == ETCP_KEEPALIVE || pkt_code == ETCP_KEEPALIVE_REQ || pkt_code == ETCP_KEEPALIVE_RESP);
standby_log_rx(is_ka ? (link->is_tcp ? "tcp keepalive" : "udp keepalive")
: (link->is_tcp ? "tcp data" : "udp data"), who);
}
#endif
if (etcp_keepalive_on_recv(e_sock, pkt, link, pkt_len)) {
return 0;
}
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) return ret;
return 0;
}
if (link->link_state == LINK_STATE_TRY_RECONNECT) {// 0 - just init, 1 - handshake, 2 - try reconnect, 3 - connected
start_keepalive_timer(link);
etcp_link_send_keepalive(link);
int old_state = link->link_state; link->link_state = LINK_STATE_CONNECTED;
etcp_fire_link_status_cbk(link, old_state, link->link_status);
}
if (link->link_state == LINK_STATE_CONNECTED) {
if (memory_pool_is_freed(link->etcp->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(link->etcp->instance->pkt_pool, pkt);
}
return 0;
}
/* Создаёт один серверный сокет из CFG_SERVER: UDP (etcp_socket_add) или TCP
* (tcp_socket_add + stcp_server_listen + stcp_server_list_add).
* Линки к соединениям не добавляет — это забота вызывающего.
* Возвращает 0 при успехе и *out_sock = созданный сокет, -1 при ошибке. */
int etcp_create_server_socket(struct UTUN_INSTANCE* instance, struct CFG_SERVER* server, struct ETCP_SOCKET** out_sock) {
if (!instance || !server || !out_sock) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_create_server_socket: invalid args"); return -1; }
*out_sock = NULL;
if (server->transport) {
uint16_t port = 0;
if (server->ip.ss_family == AF_INET) port = ntohs(((struct sockaddr_in*)&server->ip)->sin_port);
else if (server->ip.ss_family == AF_INET6) port = ntohs(((struct sockaddr_in6*)&server->ip)->sin6_port);
struct ETCP_SOCKET* ts = tcp_socket_add(instance, server);
if (!ts) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "tcp_socket_add failed for %s", server->name); return -1; }
struct stcp_link_config scfg = {.inst = instance, .listen_family = server->ip.ss_family,
.reality_enabled = server->reality_enabled};
struct stcp_server* tsrv = stcp_server_listen(&scfg, port, tcp_server_on_link, ts);
if (!tsrv) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create TCP server for %s", server->name);
tcp_socket_remove(ts);
return -1;
}
stcp_server_list_add(instance, tsrv);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server %s on port %u", server->name, port);
*out_sock = ts;
return 0;
}
struct ETCP_SOCKET* e_sock = etcp_socket_add(instance, server);
if (!e_sock) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to create socket for server %s", server->name); return -1; }
*out_sock = e_sock;
return 0;
}
// Создание всех listen-сокетов из конфига ([server] секции); автосокеты пропускаются.
// Вызывается до topo_group_init() чтобы заполнить etcp_sockets для nodeinfo.
// Возврат: 0 = OK, 1 = частичный успех (часть сокетов не создана), -1 = фатально (нет сокетов).
int init_sockets(struct UTUN_INSTANCE* instance) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
if (!instance || !instance->config) return -1;
if (instance->etcp_sockets || instance->stcp_servers) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets already initialized, skipping");
return 0;
}
struct utun_config* config = instance->config;
/* автосокеты — ручные [server] сокеты пропускаются, create_sockets вызовется отдельно */
if (config->global.auto_sockets) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "auto_sockets=enabled, skipping manual server sockets");
instance->auto_socket_enabled = 1;
return 0;
}
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);
}
}
}
}
struct ETCP_SOCKET* sock = NULL;
if (etcp_create_server_socket(instance, server, &sock) != 0) {
fail_count++;
server = server->next;
continue;
}
if (sock && !server->transport && default_ip != 0) {
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=%s ip=%s", server->name, server_type_str(server->type), ip_to_str(&addr, AF_INET).str);
}
}
if (sock && !server->transport && have_default_ip6) {
memcpy(sock->local_defaultroute_ip6, default_ip6, 16);
}
success_count++;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Initialized server %s on %s (links: %d)",
server->name, sockaddr_storage_to_str(&server->ip).str, queue_entry_count(sock->links_queue));
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
}
/* Создать TCP-линк(и) с REALITY-камуфляжем для [client] с reality=1.
* conn создаётся через NCD (только public_key, без адресов — линки добавляем сами).
* local_srv у link опционален: если задан — это [server] transport=tcp для bind
* на нужный интерфейс (ip); если нет — bind не делается (ОС выбирает source).
* Возвращает NCD-handle (сохранить в config_conn_handles) или NULL. */
struct NODE_CONN_DIRECT* etcp_config_client_reality(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client,
const uint8_t pubkey_bin[SC_PUBKEY_SIZE], uint64_t node_id) {
if (!instance || !client || !pubkey_bin) return NULL;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = client->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
struct NODE_CONN_DIRECT* handle = NULL;
int r = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, NULL);
if (r == NCD_ERR || !handle) {
DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: ncd open failed node=0x%016llx", client->name, (unsigned long long)node_id);
return NULL;
}
struct ETCP_CONN* conn = node_conn_direct_get_conn(handle);
if (!conn) { node_conn_direct_close(handle); return NULL; }
int link_count = 0;
for (struct CFG_CLIENT_LINK* cl = client->links; cl; cl = cl->next) {
if (cl->remote_addr.ss_family != AF_INET && cl->remote_addr.ss_family != AF_INET6) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: link without remote addr, skipping", client->name);
continue;
}
struct ETCP_SOCKET* bind_sock = NULL;
if (cl->local_srv) {
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (strcmp(cl->local_srv->name, s->name) == 0) { bind_sock = s; break; }
if (!bind_sock) DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: bind socket '%s' not found, using auto", client->name, cl->local_srv->name);
else if (!bind_sock->is_tcp) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: bind socket '%s' is not TCP, ignoring bind", client->name, cl->local_srv->name); bind_sock = NULL; }
}
struct ETCP_LINK* tlink = etcp_link_new(conn, bind_sock, &cl->remote_addr, 0);
if (!tlink) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: etcp_link_new failed", client->name); continue; }
tlink->is_tcp = 1;
tlink->reality = client->reality; tlink->reality_set = 1;
etcp_tcp_link_start_connect(tlink, &cl->remote_addr, 0);
link_count++;
DEBUG_INFO(DEBUG_CATEGORY_REALITY, "client %s reality: TCP link %d → %s sn=%s bind=%s",
client->name, tlink->local_link_id, sockaddr_storage_to_str(&cl->remote_addr).str,
client->reality.server_name, bind_sock ? bind_sock->name : "auto");
}
if (link_count == 0) DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: no links created", client->name);
return handle;
}
// Инициализация сети: listen-сокеты (если ещё нет) + client-соединения через NCD из конфига.
int init_connections(struct UTUN_INSTANCE* instance) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "");
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_servers) {
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;
}
}
etcp_keepalive_register(instance);
// Initialize clients via node_conn_direct
struct CFG_CLIENT* client = config->clients;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "init_connections: clients=%p total_conns=%d", config->clients, queue_entry_count(instance->connections));
while (client) {
if (strlen(client->peer_public_key_hex) == 0) {
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "no peer public key configured for client %s", client->name);
client = client->next; continue;
}
uint8_t pubkey_bin[SC_PUBKEY_SIZE];
if (sc_hex_to_binary(client->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid peer pubkey hex for client %s", client->name);
client = client->next; continue;
}
uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "client %s node_id=0x%016llx", client->name, (unsigned long long)node_id);
if (client->reality_enabled) {
struct NODE_CONN_DIRECT* rhandle = etcp_config_client_reality(instance, client, pubkey_bin, node_id);
if (rhandle) {
struct CONFIG_CONN_HANDLE* ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, client->name, MAX_CONN_NAME_LEN - 1);
ch->handle = rhandle; ch->next = instance->config_conn_handles;
instance->config_conn_handles = ch;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s saved reality handle to config_conn_handles", client->name); }
}
client = client->next; continue;
}
struct NODE_CONN_DIRECT* handle = NULL;
struct ETCP_CONN* conn = NULL;
for (struct CFG_CLIENT_LINK* cl = client->links; cl; cl = cl->next) {
if (cl->local_srv && cl->local_srv->transport) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "client %s TCP transport not yet supported via NCD, skipping", client->name);
continue;
}
struct ETCP_SOCKET* sock = NULL;
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (cl->local_srv && strcmp(cl->local_srv->name, s->name) == 0) { sock = s; break; }
if (!sock) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s link",
cl->local_srv ? cl->local_srv->name : "?", client->name);
continue;
}
if (!handle) {
struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = client->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
if (cl->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&cl->remote_addr;
v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4);
v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v4_addrs = &v4_addr;
} else if (cl->remote_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&cl->remote_addr;
v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16);
v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v6_addrs = &v6_addr;
}
int r = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock);
if (r == NCD_ERR || !handle) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd open failed for client %s", client->name);
break;
}
conn = node_conn_direct_get_conn(handle);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s ncd handle=%p conn=%p %s sock=%s",
client->name, handle, conn, r == NCD_NEW ? "NEW" : "REUSED", sock->name);
} else {
struct ETCP_LINK* link = etcp_link_new(conn, sock, &cl->remote_addr, 0);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "client %s added link link_id=%d sock=%s", client->name, link ? link->local_link_id : -1, sock->name);
}
}
if (handle) {
struct CONFIG_CONN_HANDLE* ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, client->name, MAX_CONN_NAME_LEN - 1);
ch->handle = handle; ch->next = instance->config_conn_handles;
instance->config_conn_handles = ch;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s saved handle to config_conn_handles", client->name); }
}
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;
}