Browse Source

refactor: extract handle_ping, handle_pong, send_init_response, handle_init_response_client from read callback

Reduce etcp_connections_read_callback_socket from ~600 to ~300 lines.
All 4 helpers are static file-scope functions, no logic changes.
feature/x25519-migration
Evgeny 3 months ago
parent
commit
8f13f0ec4d
  1. 595
      src/etcp_connections.c

595
src/etcp_connections.c

@ -1148,6 +1148,306 @@ int etcp_send_ping(struct UTUN_INSTANCE* instance, const uint8_t* peer_pubkey_bi
cb, user_arg, user_data, user_data_len);
}
// === Helpers extracted from etcp_connections_read_callback_socket ===
static int handle_ping(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, const struct sockaddr_storage* addr, const uint8_t* decrypted_pubkey, size_t pkt_len) {
if (pkt_len < 22) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PING too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str);
return 7;
}
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = be16toh(*(uint16_t*)(pkt->data + 17));
const uint8_t* udata = (ulen > 0) ? (pkt->data + 19) : NULL;
struct ETCP_DGRAM* resp = u_malloc(PACKET_DATA_SIZE);
if (resp) {
resp->link = NULL;
resp->noencrypt_len = SC_PUBKEY_ENC_SIZE;
uint8_t* p = resp->data;
*p++ = ETCP_PONG;
uint64_t nid = htobe64(e_sock->instance->node_id);
memcpy(p, &nid, 8); p += 8;
uint64_t n = htobe64(nonce);
memcpy(p, &n, 8); p += 8;
uint16_t resp_ulen_be = htobe16(ulen);
memcpy(p, &resp_ulen_be, 2); p += 2;
size_t copied = etcp_build_ping_response_data(udata, ulen, p, PACKET_DATA_SIZE - (p - resp->data) - SC_PUBKEY_ENC_SIZE);
p += copied;
uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE];
random_bytes(salt, sizeof(salt));
memcpy(p, salt, SC_PUBKEY_ENC_SALT_SIZE); p += SC_PUBKEY_ENC_SALT_SIZE;
uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE];
sc_obfuscate_pubkey(salt, decrypted_pubkey, e_sock->instance->my_keys.public_key, obfuscated_pubkey);
memcpy(p, obfuscated_pubkey, SC_PUBKEY_SIZE); p += SC_PUBKEY_SIZE;
resp->data_len = (uint16_t)(p - resp->data);
sc_context_t resp_sc;
sc_init_ctx(&resp_sc, &e_sock->instance->my_keys);
if (sc_set_peer_public_key(&resp_sc, decrypted_pubkey, SC_PEER_PUBKEY_BIN) == SC_OK) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG send nonce=%016llx to=%s fd=%d",
(unsigned long long)nonce, sockaddr_storage_to_str(addr).str, e_sock->fd);
etcp_send_ping_raw(resp, e_sock->fd, &resp_sc, addr);
}
u_free(resp);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
static int handle_pong(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, const struct sockaddr_storage* addr, size_t pkt_len) {
if (pkt_len < 20) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PONG too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(addr).str);
return 7;
}
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = 0;
const uint8_t* udata = NULL;
if (pkt->data_len >= 19) {
ulen = be16toh(*(uint16_t*)(pkt->data + 17));
if (ulen > 0 && pkt->data_len >= 19 + ulen) {
udata = pkt->data + 19;
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG recv nonce=%016llx data_len=%u from=%s socket=%s",
(unsigned long long)nonce, (unsigned)pkt->data_len,
sockaddr_storage_to_str(addr).str, e_sock->name);
struct PING_CONTEXT* ctx = e_sock->instance->pending_pings;
struct PING_CONTEXT* prev = NULL;
int found = 0;
while (ctx) {
if (ctx->nonce == nonce) {
found = 1;
if (prev) prev->next = ctx->next;
else e_sock->instance->pending_pings = ctx->next;
if (ctx->timeout_timer) {
uasync_cancel_timeout(e_sock->instance->ua, ctx->timeout_timer);
ctx->timeout_timer = NULL;
}
uint64_t now = get_time_tb();
uint16_t rtt = (now >= ctx->send_time) ? (uint16_t)(now - ctx->send_time) : 0;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG matched nonce=%016llx rtt=%u",
(unsigned long long)nonce, (unsigned)rtt);
ctx->cb(1, rtt, ctx->arg, nonce, udata, ulen);
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
break;
}
prev = ctx;
ctx = ctx->next;
}
if (!found) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "PONG nonce=%016llx NOT FOUND in pending (timeout?)",
(unsigned long long)nonce);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, struct ETCP_LINK* link, struct ETCP_CONN* conn, const struct ETCP_INIT_REQUEST_PKT* req, const struct sockaddr_storage* addr, uint8_t send_reset, uint32_t req_src_ip, uint16_t req_src_port, size_t pkt_len) {
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data;
// Set response code: 0x03 (with reset) or 0x05 (without reset)
// response with init (0x03) only if reinit was actually done on server side
if (send_reset != 0) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response with reset");
resp->code = ETCP_INIT_RESPONSE; // 0x03 - with reset
} else {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response without reset");
resp->code = ETCP_INIT_RESPONSE_NOINIT; // 0x05 - without reset
}
*(uint64_t*)resp->node_id = htobe64(e_sock->instance->node_id);
*(uint32_t*)resp->session_id = htobe32(conn->session_id);
resp->mtu[0]=link->mtu_local>>8;
resp->mtu[1]=link->mtu_local;
resp->link_id = link->local_link_id;
resp->remote_socket_id = req->socket_id;
resp->only_local = e_sock->only_local;
resp->type = e_sock->type;
// Add client's IP:port (so client behind NAT can know its external address)
if (addr->ss_family == AF_INET) {
struct sockaddr_in *sin = (struct sockaddr_in*)addr;
memcpy(resp->peer_ipv4, &sin->sin_addr.s_addr, 4);
uint16_t port = ntohs(sin->sin_port);
resp->peer_port[0] = port >> 8;
resp->peer_port[1] = port & 0xFF;
link->nat_ip = sin->sin_addr.s_addr;
link->nat_port = port;
} else {
// For IPv6, set to 0 (not supported for NAT traversal)
memset(resp->peer_ipv4, 0, 4);
memset(resp->peer_port, 0, 2);
link->nat_ip = 0;
link->nat_port = 0;
}
// DIRECT detection: if client reports its own address and it matches observed source → real public IP
if (pkt_len >= ETCP_INIT_REQ_V2_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->bgp) {
route_bgp_send_nat_info(link->etcp, link->remote_socket_id,
link->nat_ip, link->nat_port, NAT_TYPE_DIRECT);
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "DIRECT IP: %s:%u for %s",
ip_to_str(&link->nat_ip, AF_INET).str, link->nat_port,
link->etcp->log_name);
}
}
pkt->noencrypt_len=0;
pkt->link=link;
link->recv_keepalive = 1;
link->last_recv_local_time = get_time_tb();
link->last_recv_timestamp = pkt->timestamp;
int xoffset=sizeof(struct ETCP_INIT_RESPONSE_PKT);
// padding
int s = rand() % (link->handshake_maxsize - link->handshake_minsize) + link->handshake_minsize;
if (s > (int)(link->mtu)) s = (int)(link->mtu);
if (s < 0) s = 0;
int to_add=s - xoffset - UDP_HDR_SIZE - UDP_SC_HDR_SIZE;
if (to_add<0) to_add=0;
if (xoffset + to_add > PACKET_DATA_SIZE) { to_add = PACKET_DATA_SIZE - xoffset; if (to_add<0) to_add=0; }
for (int i=0; i<to_add; i++) pkt->data[xoffset++]=rand();// fill pad
// padding end
pkt->data_len=xoffset;
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "Sending INIT RESPONSE, link=%p, local_link_id=%d, remote_link_id=%d", link, link->local_link_id, link->remote_link_id);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP DEBUG] Send INIT RESPONSE");
etcp_encrypt_send(pkt);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
link->initialized = 1;
link->link_state = 3;
if (link->init_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (link->etcp->initialized == 0) {
etcp_conn_ready(link->etcp);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Connection established: log_name=%s socket=%s link_id=%d status=UP", link->etcp->log_name, e_sock->name, link->local_link_id);
}
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (server)", link->etcp->log_name, link, link->local_link_id);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
// Restart NAT check after link is up (e.g. after address change or reinit)
if (link->etcp->instance->bgp) {
if (link->nat_check_status < NAT_CHECK_IN_PROGRESS) route_bgp_start_link_nat_check(link->etcp->instance->bgp, link);
}
}
static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, struct ETCP_LINK* link, uint8_t pkt_code, size_t pkt_len) {
if (pkt_len < ETCP_INIT_RESP_V1_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE too short: pkt_len=%zu", pkt_len);
return 46;
}
// ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN
// ETCP_INIT_RESPONSE_NOINIT (0x05) - no reset
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data;
uint64_t server_node_id = be64toh(*(uint64_t*)resp->node_id);
uint32_t resp_session_id = be32toh(*(uint32_t*)resp->session_id);
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "INIT_RESPONSE session_id=%08x", resp_session_id);
// Check session_id: ignore response if it doesn't match our session
if (resp_session_id != link->etcp->session_id) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] INIT_RESPONSE session_id mismatch: got %08x, expected %08x, ignoring",
link->etcp->log_name, resp_session_id, link->etcp->session_id);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
link->mtu_remote = be16toh(*(uint16_t*)resp->mtu);
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU;
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote;
etcp_update_mtu(link->etcp);
link->remote_link_id = resp->link_id;
link->remote_socket_id = resp->remote_socket_id;
link->remote_only_local = resp->only_local;
link->remote_type = resp->type;
// Parse NAT IP:port from response (new format includes 4+2 bytes)
if (pkt_len >= ETCP_INIT_RESP_V2_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
{
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 {
// Legacy format without NAT info
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Received legacy INIT_RESPONSE without NAT info",
link->etcp->log_name);
}
// DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Received INIT_RESPONSE from server_node_id=%llu, mtu=%d", (unsigned long long)server_node_id, link->mtu);
link->etcp->peer_node_id = server_node_id; // If not set
etcp_update_log_name(link->etcp); // Update log_name with peer_node_id
link->initialized = 1;// получен init response (client)
link->link_state = 3; // connected
if (link->init_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (pkt_code == ETCP_INIT_RESPONSE) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit 3 %p", link->etcp);
etcp_conn_reinit(link->etcp);
}
if (link->etcp->initialized == 0) {
etcp_conn_ready(link->etcp);
}
loadbalancer_link_ready(link);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (client): rk=%d lk=%d up=%d ki=%d",
link->etcp->log_name, link, link->local_link_id, link->recv_keepalive, link->remote_keepalive, link->link_status, link->keepalive_interval);
// Start keepalive timer
etcp_link_send_keepalive(link);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp client: Link initialized successfully! Server node_id=%016llx, mtu=%d, local_link_id=%d, remote_link_id=%d", (unsigned long long)server_node_id, link->mtu, link->local_link_id, link->remote_link_id);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return 0;
}
void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connections_read_callback fd=%d, socket=%p", fd, arg);
@ -1252,94 +1552,13 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
uint8_t code = pkt->data[0];
uint64_t peer_id = be64toh(*(uint64_t*)(pkt->data + 1));
if (code == ETCP_PING) {
if (pkt_len < 22) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PING too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(&addr).str);
errorcode=7;
goto ec_fr;
}
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = be16toh(*(uint16_t*)(pkt->data + 17));
const uint8_t* udata = (ulen > 0) ? (pkt->data + 19) : NULL;
struct ETCP_DGRAM* resp = u_malloc(PACKET_DATA_SIZE);
if (resp) {
resp->link = NULL;
resp->noencrypt_len = SC_PUBKEY_ENC_SIZE;
uint8_t* p = resp->data;
*p++ = ETCP_PONG;
uint64_t nid = htobe64(e_sock->instance->node_id);
memcpy(p, &nid, 8); p += 8;
uint64_t n = htobe64(nonce);
memcpy(p, &n, 8); p += 8;
uint16_t resp_ulen_be = htobe16(ulen);
memcpy(p, &resp_ulen_be, 2); p += 2;
size_t copied = etcp_build_ping_response_data(udata, ulen, p, PACKET_DATA_SIZE - (p - resp->data) - SC_PUBKEY_ENC_SIZE);
p += copied;
uint8_t salt[SC_PUBKEY_ENC_SALT_SIZE];
random_bytes(salt, sizeof(salt));
memcpy(p, salt, SC_PUBKEY_ENC_SALT_SIZE); p += SC_PUBKEY_ENC_SALT_SIZE;
uint8_t obfuscated_pubkey[SC_PUBKEY_SIZE];
sc_obfuscate_pubkey(salt, decrypted_pubkey, e_sock->instance->my_keys.public_key, obfuscated_pubkey);
memcpy(p, obfuscated_pubkey, SC_PUBKEY_SIZE); p += SC_PUBKEY_SIZE;
resp->data_len = (uint16_t)(p - resp->data);
sc_context_t resp_sc;
sc_init_ctx(&resp_sc, &e_sock->instance->my_keys);
if (sc_set_peer_public_key(&resp_sc, decrypted_pubkey, SC_PEER_PUBKEY_BIN) == SC_OK) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG send nonce=%016llx to=%s fd=%d",
(unsigned long long)nonce, sockaddr_storage_to_str(&addr).str, e_sock->fd);
etcp_send_ping_raw(resp, e_sock->fd, &resp_sc, &addr);
}
u_free(resp);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
int ret = handle_ping(e_sock, pkt, &addr, decrypted_pubkey, pkt_len);
if (ret) { errorcode = ret; goto ec_fr; }
return;
}
if (code == ETCP_PONG) {
if (pkt_len < 20) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "PONG too short: pkt_len=%zu from %s", pkt_len, sockaddr_storage_to_str(&addr).str);
errorcode=7;
goto ec_fr;
}
uint64_t nonce = be64toh(*(uint64_t*)(pkt->data + 9));
uint16_t ulen = 0;
const uint8_t* udata = NULL;
if (pkt->data_len >= 19) {
ulen = be16toh(*(uint16_t*)(pkt->data + 17));
if (ulen > 0 && pkt->data_len >= 19 + ulen) {
udata = pkt->data + 19;
}
}
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG recv nonce=%016llx data_len=%u from=%s socket=%s",
(unsigned long long)nonce, (unsigned)pkt->data_len,
sockaddr_storage_to_str(&addr).str, e_sock->name);
struct PING_CONTEXT* ctx = e_sock->instance->pending_pings;
struct PING_CONTEXT* prev = NULL;
int found = 0;
while (ctx) {
if (ctx->nonce == nonce) {
found = 1;
if (prev) prev->next = ctx->next;
else e_sock->instance->pending_pings = ctx->next;
if (ctx->timeout_timer) {
uasync_cancel_timeout(e_sock->instance->ua, ctx->timeout_timer);
ctx->timeout_timer = NULL;
}
uint64_t now = get_time_tb();
uint16_t rtt = (now >= ctx->send_time) ? (uint16_t)(now - ctx->send_time) : 0;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "PONG matched nonce=%016llx rtt=%u",
(unsigned long long)nonce, (unsigned)rtt);
ctx->cb(1, rtt, ctx->arg, nonce, udata, ulen);
if (ctx->user_data) u_free(ctx->user_data);
u_free(ctx);
break;
}
prev = ctx;
ctx = ctx->next;
}
if (!found) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "PONG nonce=%016llx NOT FOUND in pending (timeout?)",
(unsigned long long)nonce);
}
memory_pool_free(e_sock->instance->pkt_pool, pkt);
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) {
@ -1488,100 +1707,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
uint32_t req_src_ip; memcpy(&req_src_ip, req->src_ipv4, 4);
uint16_t req_src_port = be16toh(*(uint16_t*)req->src_port);
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data;
// Set response code: 0x03 (with reset) or 0x05 (without reset)
// response with init (0x03) only if reinit was actually done on server side
if (send_reset != 0) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response with reset");
resp->code = ETCP_INIT_RESPONSE; // 0x03 - with reset
} else {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "send init_response without reset");
resp->code = ETCP_INIT_RESPONSE_NOINIT; // 0x05 - without reset
}
*(uint64_t*)resp->node_id = htobe64(e_sock->instance->node_id);
*(uint32_t*)resp->session_id = htobe32(conn->session_id);
resp->mtu[0]=link->mtu_local>>8;
resp->mtu[1]=link->mtu_local;
resp->link_id = link->local_link_id;
resp->remote_socket_id = req->socket_id;
resp->only_local = e_sock->only_local;
resp->type = e_sock->type;
// Add client's IP:port (so client behind NAT can know its external address)
if (addr.ss_family == AF_INET) {
struct sockaddr_in *sin = (struct sockaddr_in*)&addr;
memcpy(resp->peer_ipv4, &sin->sin_addr.s_addr, 4);
uint16_t port = ntohs(sin->sin_port);
resp->peer_port[0] = port >> 8;
resp->peer_port[1] = port & 0xFF;
link->nat_ip = sin->sin_addr.s_addr;
link->nat_port = port;
} else {
// For IPv6, set to 0 (not supported for NAT traversal)
memset(resp->peer_ipv4, 0, 4);
memset(resp->peer_port, 0, 2);
link->nat_ip = 0;
link->nat_port = 0;
}
// DIRECT detection: if client reports its own address and it matches observed source → real public IP
if (pkt_len >= ETCP_INIT_REQ_V2_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->bgp) {
route_bgp_send_nat_info(link->etcp, link->remote_socket_id,
link->nat_ip, link->nat_port, NAT_TYPE_DIRECT);
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "DIRECT IP: %s:%u for %s",
ip_to_str(&link->nat_ip, AF_INET).str, link->nat_port,
link->etcp->log_name);
}
}
pkt->noencrypt_len=0;
pkt->link=link;
link->recv_keepalive = 1;
link->last_recv_local_time = get_time_tb();
link->last_recv_timestamp = pkt->timestamp;
int xoffset=sizeof(struct ETCP_INIT_RESPONSE_PKT);
// padding
int s = rand() % (link->handshake_maxsize - link->handshake_minsize) + link->handshake_minsize;
if (s > (int)(link->mtu)) s = (int)(link->mtu);
if (s < 0) s = 0;
int to_add=s - xoffset - UDP_HDR_SIZE - UDP_SC_HDR_SIZE;
if (to_add<0) to_add=0;
if (xoffset + to_add > PACKET_DATA_SIZE) { to_add = PACKET_DATA_SIZE - xoffset; if (to_add<0) to_add=0; }
for (int i=0; i<to_add; i++) pkt->data[xoffset++]=rand();// fill pad
// padding end
pkt->data_len=xoffset;
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "Sending INIT RESPONSE, link=%p, local_link_id=%d, remote_link_id=%d", link, link->local_link_id, link->remote_link_id);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP DEBUG] Send INIT RESPONSE");
etcp_encrypt_send(pkt);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
link->initialized = 1;
link->link_state = 3;
if (link->init_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (link->etcp->initialized == 0) {
etcp_conn_ready(link->etcp);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "Connection established: log_name=%s socket=%s link_id=%d status=UP", link->etcp->log_name, e_sock->name, link->local_link_id);
}
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (server)", link->etcp->log_name, link, link->local_link_id);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
// Restart NAT check after link is up (e.g. after address change or reinit)
if (link->etcp->instance->bgp) {
if (link->nat_check_status < NAT_CHECK_IN_PROGRESS) route_bgp_start_link_nat_check(link->etcp->instance->bgp, link);
}
send_init_response(e_sock, pkt, link, conn, req, &addr, send_reset, req_src_ip, req_src_port, pkt_len);
return;
@ -1627,114 +1753,9 @@ process_decrypted:
}
if (pkt_code == ETCP_INIT_RESPONSE || pkt_code == ETCP_INIT_RESPONSE_NOINIT) {
if (pkt_len < ETCP_INIT_RESP_V1_SIZE) { errorcode = 46; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "INIT_RESPONSE too short: pkt_len=%zu", pkt_len); goto ec_fr; }
// ETCP_INIT_RESPONSE (0x03) - reset entire ETCP_CONN
// ETCP_INIT_RESPONSE_NOINIT (0x05) - no reset
// Save src_ipv4/src_port from INIT REQUEST before we overwrite the buffer with INIT RESPONSE
uint32_t req_src_ip; memcpy(&req_src_ip, req->src_ipv4, 4);
uint16_t req_src_port = be16toh(*(uint16_t*)req->src_port);
struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data;
uint64_t server_node_id = be64toh(*(uint64_t*)resp->node_id);
uint32_t resp_session_id = be32toh(*(uint32_t*)resp->session_id);
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "INIT_RESPONSE session_id=%08x", resp_session_id);
// Check session_id: ignore response if it doesn't match our session
if (resp_session_id != link->etcp->session_id) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] INIT_RESPONSE session_id mismatch: got %08x, expected %08x, ignoring",
link->etcp->log_name, resp_session_id, link->etcp->session_id);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return;
}
link->mtu_remote = be16toh(*(uint16_t*)resp->mtu);
if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU;
link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote;
etcp_update_mtu(link->etcp);
link->remote_link_id = resp->link_id;
link->remote_socket_id = resp->remote_socket_id;
link->remote_only_local = resp->only_local;
link->remote_type = resp->type;
// Parse NAT IP:port from response (new format includes 4+2 bytes)
if (pkt_len >= ETCP_INIT_RESP_V2_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
{
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 {
// Legacy format without NAT info
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] Received legacy INIT_RESPONSE without NAT info",
link->etcp->log_name);
}
// DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Received INIT_RESPONSE from server_node_id=%llu, mtu=%d", (unsigned long long)server_node_id, link->mtu);
link->etcp->peer_node_id = server_node_id; // If not set
etcp_update_log_name(link->etcp); // Update log_name with peer_node_id
link->initialized = 1;// получен init response (client)
link->link_state = 3; // connected
if (link->init_timer) {
uasync_cancel_timeout(link->etcp->instance->ua, link->init_timer);
link->init_timer = NULL;
}
if (pkt_code == ETCP_INIT_RESPONSE) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit 3 %p", link->etcp);
etcp_conn_reinit(link->etcp);
}
if (link->etcp->initialized == 0) {
etcp_conn_ready(link->etcp);
}
loadbalancer_link_ready(link);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %p (local_id=%d) initialized and marked as UP (client): rk=%d lk=%d up=%d ki=%d",
link->etcp->log_name, link, link->local_link_id, link->recv_keepalive, link->remote_keepalive, link->link_status, link->keepalive_interval);
// Start keepalive timer
etcp_link_send_keepalive(link);
start_keepalive_timer(link);
loadbalancer_link_ready(link);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "etcp client: Link initialized successfully! Server node_id=%016llx, mtu=%d, local_link_id=%d, remote_link_id=%d", (unsigned long long)server_node_id, link->mtu, link->local_link_id, link->remote_link_id);
memory_pool_free(e_sock->instance->pkt_pool, pkt);
return; // INIT_RESPONSE is handled, no further processing needed
int ret = handle_init_response_client(e_sock, pkt, link, pkt_code, pkt_len);
if (ret) { errorcode = ret; goto ec_fr; }
return;
}
if (link->link_state == 2) {// из recovery получен нормальный пакет - восстанавливаем линк в нормальный режим

Loading…
Cancel
Save