Browse Source

fix: uasync epoll event ordering + test_uasync_socket_race fix

topo_upd
evgeny 2 months ago
parent
commit
ce5a84da24
  1. 52
      lib/u_async.c
  2. 195
      src/routing_layer/conn_mgr.c
  3. 61
      src/transport_layer/etcp.c
  4. 1
      src/transport_layer/etcp.h
  5. 166
      src/transport_layer/etcp_connect.c
  6. 156
      src/transport_layer/etcp_connect.h
  7. 6
      src/transport_layer/etcp_connections.c
  8. 8
      tests/test_uasync_socket_race.c

52
lib/u_async.c

@ -941,14 +941,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
socket_t_callback_t local_read_sock = node->read_cbk_sock; socket_t_callback_t local_read_sock = node->read_cbk_sock;
socket_t_callback_t local_write_sock = node->write_cbk_sock; socket_t_callback_t local_write_sock = node->write_cbk_sock;
/* Check for error conditions first */ /* Read readiness BEFORE error — avoid losing data on combined IN+HUP events */
if (events[i].events & (EPOLLERR | EPOLLHUP)) {
if (local_except) {
local_except(local_fd, local_ud);
}
}
/* Read readiness - use appropriate callback based on socket type */
if (events[i].events & EPOLLIN) { if (events[i].events & EPOLLIN) {
if (local_type == SOCKET_NODE_TYPE_SOCK) { if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_read_sock) { if (local_read_sock) {
@ -961,7 +954,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
} }
} }
/* Write readiness - use appropriate callback based on socket type */ /* Write readiness BEFORE error — flush pending writes before handling HUP */
if (events[i].events & EPOLLOUT) { if (events[i].events & EPOLLOUT) {
if (local_type == SOCKET_NODE_TYPE_SOCK) { if (local_type == SOCKET_NODE_TYPE_SOCK) {
if (local_write_sock) { if (local_write_sock) {
@ -973,6 +966,13 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events,
} }
} }
} }
/* Check for error conditions LAST — I/O handlers drain/process data first */
if (events[i].events & (EPOLLERR | EPOLLHUP)) {
if (local_except) {
local_except(local_fd, local_ud);
}
}
} }
} }
#endif #endif
@ -1238,22 +1238,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
} }
if (!node) continue; // Socket may have been removed if (!node) continue; // Socket may have been removed
/* Check for error conditions first */ /* Read readiness BEFORE error — avoid losing data on combined IN+ERR/HUP events */
if (ua->poll_fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) {
/* Treat as exceptional condition */
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
/* Exceptional data (out-of-band) */
if (ua->poll_fds[i].revents & POLLPRI) {
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
/* Read readiness - use appropriate callback based on socket type */
if (ua->poll_fds[i].revents & POLLIN) { if (ua->poll_fds[i].revents & POLLIN) {
if (node->type == SOCKET_NODE_TYPE_SOCK) { if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->read_cbk_sock) { if (node->read_cbk_sock) {
@ -1266,7 +1251,7 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
} }
} }
/* Write readiness - use appropriate callback based on socket type */ /* Write readiness BEFORE error — flush pending writes before handling HUP */
if (ua->poll_fds[i].revents & POLLOUT) { if (ua->poll_fds[i].revents & POLLOUT) {
if (node->type == SOCKET_NODE_TYPE_SOCK) { if (node->type == SOCKET_NODE_TYPE_SOCK) {
if (node->write_cbk_sock) { if (node->write_cbk_sock) {
@ -1278,6 +1263,21 @@ void uasync_poll(struct UASYNC* ua, int timeout_tb) {
} }
} }
} }
/* Check for error conditions LAST — I/O handlers drain/process data first */
if (ua->poll_fds[i].revents & (POLLERR | POLLHUP | POLLNVAL)) {
/* Treat as exceptional condition */
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
/* Exceptional data (out-of-band) */
if (ua->poll_fds[i].revents & POLLPRI) {
if (node->except_cbk) {
node->except_cbk(node->fd, node->user_data);
}
}
} }
} }
#endif #endif

195
src/routing_layer/conn_mgr.c

@ -289,13 +289,6 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_
cm_deliver_result(entry, CONN_MGR_OK); cm_deliver_result(entry, CONN_MGR_OK);
return CONN_MGR_OK; return CONN_MGR_OK;
} }
entry = cm_ensure_entry(mgr, node_id);
if (!entry) return CONN_MGR_ERR_INTERNAL;
entry->state = CONN_MGR_STATE_CONNECTING;
{ struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node));
if (cn) { cn->cb = cb; cn->arg = cb_arg; cn->next = entry->cb_list; entry->cb_list = cn; } }
topo_group_new_conn(group, existing);
return CONN_MGR_OK;
} }
entry = cm_ensure_entry(mgr, node_id); entry = cm_ensure_entry(mgr, node_id);
if (!entry) return CONN_MGR_ERR_INTERNAL; if (!entry) return CONN_MGR_ERR_INTERNAL;
@ -481,6 +474,25 @@ static int cm_is_rtt_fresh(struct TOPO_NODEQ* nq, uint64_t now_tb) {
return (last > 0 && (now_tb - last) < (uint64_t)(CONN_MGR_CANDIDATE_CACHE_MS * 10)); return (last > 0 && (now_tb - last) < (uint64_t)(CONN_MGR_CANDIDATE_CACHE_MS * 10));
} }
#define CM_V6_LINK_LOCAL 0
#define CM_V6_LOCAL 1
#define CM_V6_DIRECT 2
#define CM_V6_OTHER 3
static uint8_t cm_classify_v6_addr(const uint8_t addr[16]) {
if (addr[0] == 0xfe && (addr[1] & 0xc0) == 0x80) return CM_V6_LINK_LOCAL;
if (addr[0] == 0xfc || addr[0] == 0xfd) return CM_V6_LOCAL;
if (addr[0] == 0xff) return CM_V6_OTHER;
{ uint8_t zero[16] = {0}; if (memcmp(addr, zero, 16) == 0) return CM_V6_OTHER; }
{ uint8_t lb[16] = {0,0,0,0,0,0,0,0,0,0,0,0,0,0,0,1}; if (memcmp(addr, lb, 16) == 0) return CM_V6_OTHER; }
return CM_V6_DIRECT;
}
static uint8_t cm_sock_v6_classify(const struct ETCP_SOCKET* s) {
if (s->local_addr.ss_family != AF_INET6) return CM_V6_OTHER;
return cm_classify_v6_addr(((const struct sockaddr_in6*)&s->local_addr)->sin6_addr.s6_addr);
}
static int cm_has_direct_ip(struct TOPO_NODEQ* nq) { static int cm_has_direct_ip(struct TOPO_NODEQ* nq) {
if (!nq || !nq->node) return 0; if (!nq || !nq->node) return 0;
for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) { for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) {
@ -492,6 +504,13 @@ static int cm_has_direct_ip(struct TOPO_NODEQ* nq) {
} }
} }
} }
for (const struct TOPO_ADDR6* a6 = nq->node->v6_addrs; a6; a6 = a6->next) {
if ((a6->type != TOPO_ADDR_NAT && a6->type != TOPO_ADDR_INTERFACE) || cm_classify_v6_addr(a6->addr) != CM_V6_DIRECT) continue;
for (const struct TOPO_SOCKMETA6* m = nq->node->v6_sock_meta; m; m = m->next)
if (m->id == a6->socket_id && (m->config_type == CFG_SERVER_TYPE_PUBLIC || m->config_type == CFG_SERVER_TYPE_UNKNOWN
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT))
return 1;
}
return 0; return 0;
} }
@ -505,6 +524,14 @@ static int cm_has_local_addr(struct TOPO_NODEQ* nq) {
} }
} }
} }
for (const struct TOPO_ADDR6* a6 = nq->node->v6_addrs; a6; a6 = a6->next) {
if (a6->type != TOPO_ADDR_INTERFACE) continue;
uint8_t cls = cm_classify_v6_addr(a6->addr);
if (cls != CM_V6_LINK_LOCAL && cls != CM_V6_LOCAL) continue;
for (const struct TOPO_SOCKMETA6* m = nq->node->v6_sock_meta; m; m = m->next)
if (m->id == a6->socket_id && (m->config_type == CFG_SERVER_TYPE_PRIVATE || m->config_type == CFG_SERVER_TYPE_LOCAL))
return 1;
}
return 0; return 0;
} }
@ -546,41 +573,71 @@ static void cm_start_local_scan(struct CONN_MGR_ENTRY* entry) {
entry->local_scan_state = CM_TRY_FAILED; entry->local_scan_state = CM_TRY_FAILED;
} }
static struct ETCP_CONN* cm_get_or_create_conn(struct CONN_MGR_ENTRY* entry, const struct TOPO_NODEQ* target) {
struct ETCP_CONN* conn = topo_group_find_conn_for_node(entry->mgr->group, entry->node_id);
if (!conn) conn = instance_find_conn(entry->mgr->instance, entry->node_id);
if (conn) {
if (conn->peer_node_id == entry->node_id) { struct etcp_cbk_entry* cbe = conn->init_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } conn->init_cbks = NULL; etcp_conn_add_init_cbk(conn, cm_direct_init_cb, entry); return conn; }
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: existing conn peer_id mismatch 0x%016llx != 0x%016llx, ignoring",
(unsigned long long)conn->peer_node_id, (unsigned long long)entry->node_id);
}
conn = etcp_connection_create(entry->mgr->instance, NULL);
if (!conn) return NULL;
sc_init_ctx(&conn->crypto_ctx, &entry->mgr->instance->my_keys);
sc_set_peer_public_key(&conn->crypto_ctx, target->node->public_key, 0);
conn->peer_node_id = entry->node_id;
etcp_update_log_name(conn);
etcp_conn_add_init_cbk(conn, cm_direct_init_cb, entry);
return conn;
}
static void cm_add_v6_link(struct ETCP_CONN* conn, const uint8_t addr[16], uint16_t port, struct ETCP_SOCKET* s) {
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, addr, 16); sin6.sin6_port = htons(port);
if (cm_classify_v6_addr(addr) == CM_V6_LINK_LOCAL) { struct sockaddr_in6* our = (struct sockaddr_in6*)&s->local_addr; sin6.sin6_scope_id = our->sin6_scope_id; }
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
etcp_link_new(conn, s, &sa, 0);
}
static void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) { static void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
if (entry->main_connect_state == CM_TRY_OK) return; if (entry->main_connect_state == CM_TRY_OK || entry->local_scan_state == CM_TRY_OK) return;
if (entry->local_scan_state == CM_TRY_OK) return;
struct TOPO_GROUP* group = entry->mgr->group; struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_NODEQ* target = topo_node_find_by_id(group, entry->node_id); struct TOPO_NODEQ* target = topo_node_find_by_id(group, entry->node_id);
if (!target || !target->node) { cm_deliver_result(entry, CONN_MGR_ERR_NOT_FOUND); return; } if (!target || !target->node) { cm_deliver_result(entry, CONN_MGR_ERR_NOT_FOUND); return; }
struct ETCP_CONN* conn = cm_get_or_create_conn(entry, target);
if (!conn) { cm_deliver_result(entry, CONN_MGR_ERR_INTERNAL); return; }
int any_link = 0;
const char* timer_name = entry->db_loaded ? "conn_mgr_db_direct" : "conn_mgr_direct";
if (entry->db_loaded) { if (entry->db_loaded) {
for (const struct TOPO_ADDR4* a = target->node->v4_addrs; a; a = a->next) { for (const struct TOPO_ADDR4* a = target->node->v4_addrs; a; a = a->next) {
if (a->protocol != TOPO_PROTO_UDP) continue; if (!(a->protocol & TOPO_PROTO_UDP)) continue;
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets; struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { while (s) {
if (s->local_addr.ss_family == AF_INET) { if (s->local_addr.ss_family == AF_INET) {
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
struct ETCP_CONN* conn = etcp_connection_create(entry->mgr->instance, NULL); if (etcp_link_new(conn, s, &sa, 0)) any_link = 1;
if (conn) {
etcp_conn_add_init_cbk(conn, cm_direct_init_cb, entry);
sc_init_ctx(&conn->crypto_ctx, &entry->mgr->instance->my_keys);
sc_set_peer_public_key(&conn->crypto_ctx, target->node->public_key, 0);
etcp_conn_set_peer_node_id(conn, entry->node_id);
if (etcp_link_new(conn, s, &sa, 0)) {
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, (int)(entry->mgr->direct_timeout_ms * 10), entry, cm_direct_timeout_cb, "conn_mgr_db_direct");
return;
}
etcp_connection_close(conn);
}
} }
s = s->next; s = s->next;
} }
} }
for (const struct TOPO_ADDR6* a6 = target->node->v6_addrs; a6; a6 = a6->next) {
if (!(a6->protocol & TOPO_PROTO_UDP)) continue;
uint8_t tclass = cm_classify_v6_addr(a6->addr); if (tclass == CM_V6_OTHER) continue;
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (cm_sock_v6_classify(s) == tclass) { cm_add_v6_link(conn, a6->addr, a6->port, s); any_link = 1; } s = s->next; }
}
if (!any_link) {
etcp_connection_close(conn);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: db_node 0x%016llx direct failed (no compatible addr/socket)", (unsigned long long)entry->node_id); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: db_node 0x%016llx direct failed (no compatible addr/socket)", (unsigned long long)entry->node_id);
cm_cleanup_db_node(entry); cm_cleanup_db_node(entry); cm_deliver_result(entry, CONN_MGR_ERR_UNREACHABLE);
cm_deliver_result(entry, CONN_MGR_ERR_UNREACHABLE); } else {
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, (int)(entry->mgr->direct_timeout_ms * 10), entry, cm_direct_timeout_cb, timer_name);
}
return; return;
} }
@ -595,25 +652,26 @@ static void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
if (cm_nat_compatible(s, m->config_type, m->nat_type) && s->local_addr.ss_family == AF_INET) { if (cm_nat_compatible(s, m->config_type, m->nat_type) && s->local_addr.ss_family == AF_INET) {
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin));
struct ETCP_CONN* conn = etcp_connection_create(entry->mgr->instance, NULL); if (etcp_link_new(conn, s, &sa, 0)) any_link = 1;
if (conn) {
etcp_conn_add_init_cbk(conn, cm_direct_init_cb, entry);
sc_init_ctx(&conn->crypto_ctx, &entry->mgr->instance->my_keys);
sc_set_peer_public_key(&conn->crypto_ctx, target->node->public_key, 0);
etcp_conn_set_peer_node_id(conn, entry->node_id);
if (etcp_link_new(conn, s, &sa, 0)) {
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, (int)(entry->mgr->direct_timeout_ms * 10), entry, cm_direct_timeout_cb, "conn_mgr_direct");
return;
} }
etcp_connection_close(conn); s = s->next;
} }
} }
s = s->next;
} }
for (const struct TOPO_ADDR6* a6 = target->node->v6_addrs; a6; a6 = a6->next) {
if (a6->type != priority_order[pri]) continue;
uint8_t tclass = cm_classify_v6_addr(a6->addr); if (tclass == CM_V6_OTHER) continue;
struct ETCP_SOCKET* s = entry->mgr->instance->etcp_sockets;
while (s) { if (cm_sock_v6_classify(s) == tclass) { cm_add_v6_link(conn, a6->addr, a6->port, s); any_link = 1; } s = s->next; }
} }
} }
if (any_link) {
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, (int)(entry->mgr->direct_timeout_ms * 10), entry, cm_direct_timeout_cb, timer_name);
return;
} }
etcp_connection_close(conn);
int has_our = cm_has_direct_ip(group->local_node); int has_our = cm_has_direct_ip(group->local_node);
int has_target = cm_has_direct_ip(target); int has_target = cm_has_direct_ip(target);
if (has_our && !has_target) { cm_start_phase_reverse(entry); return; } if (has_our && !has_target) { cm_start_phase_reverse(entry); return; }
@ -1102,6 +1160,7 @@ static void cm_bg_ping_timer_cb(void* arg) {
while (l) { if (l->initialized && l->link_status) { has_active = 1; break; } l = l->next; } while (l) { if (l->initialized && l->link_status) { has_active = 1; break; } l = l->next; }
} }
if (!has_active) { if (!has_active) {
int sent_v4 = 0;
for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) { for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue; if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
@ -1110,12 +1169,26 @@ static void cm_bg_ping_timer_cb(void* arg) {
struct ETCP_SOCKET* s = mgr->instance->etcp_sockets; struct ETCP_SOCKET* s = mgr->instance->etcp_sockets;
while (s) { if (s->local_addr.ss_family == AF_INET) break; s = s->next; } while (s) { if (s->local_addr.ss_family == AF_INET) break; s = s->next; }
if (!s) break; if (!s) break;
etcp_send_ping_to_socket(mgr->instance, s, nq->node->public_key, &sa,
CONN_PROBE_TIMEOUT_MS, cm_bg_ping_noop_cb, NULL, NULL, 0, 0);
sent_v4 = 1; break;
}
if (!sent_v4) {
for (const struct TOPO_ADDR6* a6 = nq->node->v6_addrs; a6; a6 = a6->next) {
if (a6->type != TOPO_ADDR_NAT && a6->type != TOPO_ADDR_INTERFACE) continue;
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, a6->addr, 16); sin6.sin6_port = htons(a6->port);
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
struct ETCP_SOCKET* s = mgr->instance->etcp_sockets;
while (s) { if (s->local_addr.ss_family == AF_INET6) break; s = s->next; }
if (!s) break;
etcp_send_ping_to_socket(mgr->instance, s, nq->node->public_key, &sa, etcp_send_ping_to_socket(mgr->instance, s, nq->node->public_key, &sa,
CONN_PROBE_TIMEOUT_MS, cm_bg_ping_noop_cb, NULL, NULL, 0, 0); CONN_PROBE_TIMEOUT_MS, cm_bg_ping_noop_cb, NULL, NULL, 0, 0);
break; break;
} }
} }
} }
}
mgr->bg_ping_cursor++; mgr->bg_ping_cursor++;
} }
mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping");
@ -1364,7 +1437,7 @@ static void cm_invite_cleanup(struct cm_invite_pending* inv) {
int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni, int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
uint64_t group_id, uint32_t timeout_ms, uint64_t group_id, uint32_t timeout_ms,
conn_mgr_connect_callback_t cb, void* cb_arg) { conn_mgr_connect_callback_t cb, void* cb_arg) {
if (!mgr || !ni || !ni->v4_addrs) { if (cb) cb(CONN_MGR_ERR_INTERNAL, ni ? ni->node_id : 0, cb_arg); return CONN_MGR_ERR_INTERNAL; } if (!mgr || !ni || (!ni->v4_addrs && !ni->v6_addrs)) { if (cb) cb(CONN_MGR_ERR_INTERNAL, ni ? ni->node_id : 0, cb_arg); return CONN_MGR_ERR_INTERNAL; }
uint64_t node_id = ni->node_id; uint64_t node_id = ni->node_id;
if (node_id == 0 || node_id == mgr->instance->node_id) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } if (node_id == 0 || node_id == mgr->instance->node_id) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
struct TOPO_GROUP* group = mgr->group; struct TOPO_GROUP* group = mgr->group;
@ -1386,7 +1459,14 @@ int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
a4->type = TOPO_ADDR_NAT; a4->protocol = src->protocol; a4->socket_id = 0; a4->type = TOPO_ADDR_NAT; a4->protocol = src->protocol; a4->socket_id = 0;
a4->next = group_ni->v4_addrs; group_ni->v4_addrs = a4; a4->next = group_ni->v4_addrs; group_ni->v4_addrs = a4;
} }
if (!group_ni->v4_addrs) { u_free(group_ni); if (cb) cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, cb_arg); return CONN_MGR_ERR_NO_ADDRESSES; } for (const struct TOPO_ADDR6* src6 = ni->v6_addrs; src6; src6 = src6->next) {
struct TOPO_ADDR6* a6 = memory_pool_alloc(groups->v6_addr_pool);
if (!a6) continue;
memset(a6, 0, sizeof(*a6)); memcpy(a6->addr, src6->addr, 16); a6->port = src6->port;
a6->type = TOPO_ADDR_NAT; a6->protocol = src6->protocol; a6->socket_id = 0;
a6->next = group_ni->v6_addrs; group_ni->v6_addrs = a6;
}
if (!group_ni->v4_addrs && !group_ni->v6_addrs) { u_free(group_ni); if (cb) cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, cb_arg); return CONN_MGR_ERR_NO_ADDRESSES; }
group_ni = topo_node_registry_acquire(groups, group_ni); /* group_ni now owned by registry (ref=1) */ group_ni = topo_node_registry_acquire(groups, group_ni); /* group_ni now owned by registry (ref=1) */
if (!group_ni) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } if (!group_ni) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
@ -1433,6 +1513,25 @@ int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
s = s->next; s = s->next;
} }
} }
for (const struct TOPO_ADDR6* a6 = group_ni->v6_addrs; a6; a6 = a6->next) {
if (!(a6->protocol & TOPO_PROTO_UDP)) continue;
uint8_t tclass = cm_classify_v6_addr(a6->addr); if (tclass == CM_V6_OTHER) continue;
struct ETCP_SOCKET* s = mgr->instance->etcp_sockets;
while (s) {
if (cm_sock_v6_classify(s) == tclass) {
if (!udp_conn) {
udp_conn = etcp_connection_create(mgr->instance, NULL);
if (udp_conn) {
etcp_conn_add_init_cbk(udp_conn, cm_invite_init_cb, inv);
sc_init_ctx(&udp_conn->crypto_ctx, &mgr->instance->my_keys);
sc_set_peer_public_key(&udp_conn->crypto_ctx, group_ni->public_key, 0);
}
}
if (udp_conn) { cm_add_v6_link(udp_conn, a6->addr, a6->port, s); inv->conn = udp_conn; udp_connected = 1; }
}
s = s->next;
}
}
/* ── TCP: first address ── */ /* ── TCP: first address ── */
for (const struct TOPO_ADDR4* a = group_ni->v4_addrs; a; a = a->next) { for (const struct TOPO_ADDR4* a = group_ni->v4_addrs; a; a = a->next) {
@ -1453,6 +1552,24 @@ int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
} }
break; break;
} }
for (const struct TOPO_ADDR6* a6 = group_ni->v6_addrs; a6; a6 = a6->next) {
if (!(a6->protocol & TOPO_PROTO_TCP)) continue;
struct sockaddr_in6 sin6; memset(&sin6, 0, sizeof(sin6)); sin6.sin6_family = AF_INET6;
memcpy(&sin6.sin6_addr, a6->addr, 16); sin6.sin6_port = htons(a6->port);
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
struct stcp_link_config tcp_cfg6 = {
.ua = mgr->instance->ua, .my_keys = &mgr->instance->my_keys, .inst = mgr->instance,
.peer_pubkey = group_ni->public_key, .peer_pubkey_mode = 0,
.remote_addr = &sa, .remote_port = a6->port,
};
inv->tcp_link = stcp_link_connect(&tcp_cfg6);
if (inv->tcp_link) {
stcp_link_set_on_ready(inv->tcp_link, cm_tcp_ready_cb, inv);
stcp_link_set_on_close(inv->tcp_link, cm_tcp_close_cb, inv);
tcp_connected = 1;
}
break;
}
if (!udp_connected && !tcp_connected) { if (!udp_connected && !tcp_connected) {
cm_invite_fail(inv, CONN_MGR_ERR_UNREACHABLE); return CONN_MGR_ERR_UNREACHABLE; cm_invite_fail(inv, CONN_MGR_ERR_UNREACHABLE); return CONN_MGR_ERR_UNREACHABLE;

61
src/transport_layer/etcp.c

@ -652,68 +652,7 @@ void etcp_conn_queue_set_ready(struct ETCP_CONN* conn) {
conn->callbacks_running = 0; conn->callbacks_running = 0;
} }
// Set peer_node_id, reindex if already in connections queue
void etcp_conn_set_peer_node_id(struct ETCP_CONN* conn, uint64_t peer_node_id) {
if (!conn) return;
if (conn->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] set_peer_node_id on deleted conn", conn->log_name); return; }
if (conn->peer_node_id == peer_node_id) return;
if (peer_node_id != 0 && conn->instance && conn->instance->connections) {
struct ll_entry* clash = queue_find_data_by_index(conn->instance->connections, (const uint8_t*)&peer_node_id);
while (clash) {
struct conn_queue_entry* cqe = (struct conn_queue_entry*)clash->data;
if (cqe && cqe->conn && cqe->conn != conn && cqe->peer_node_id == peer_node_id) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP,
"!!!!!!!!!!!! PEER NODE ID COLLISION !!!!!!!!!!!! "
"conn=%s trying peer=0x%016llx already_owned_by=%s(peer=0x%016llx state=%d)",
conn->log_name, (unsigned long long)peer_node_id,
cqe->conn->log_name, (unsigned long long)cqe->conn->peer_node_id, cqe->conn->state);
return;
}
clash = queue_find_next_by_index(conn->instance->connections, (const uint8_t*)&peer_node_id, clash);
}
}
if (conn->state == 0 && conn->peer_node_id == 0 && peer_node_id != 0) {
/* первый раз ставим реальный peer_id: переиндексировать из key=0 в key=peer_node_id */
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id 0 -> 0x%llx, reindexing pending conn",
conn->log_name, (unsigned long long)peer_node_id);
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] reindex alloc failed", conn->log_name); return; }
struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = peer_node_id;
ce->conn = conn;
conn->conn_queue_entry = qe;
conn->conn_queue = conn->instance->connections;
queue_data_put_with_index(conn->instance->connections, qe);
} else if (conn->state == 1 && conn->peer_node_id != 0) {
// Already in indexed queue with old ID — reindex
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id changed 0x%llx -> 0x%llx, reindexing",
conn->log_name, (unsigned long long)conn->peer_node_id, (unsigned long long)peer_node_id);
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] reindex alloc failed", conn->log_name); return; }
struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data;
ce->peer_node_id = peer_node_id;
ce->conn = conn;
conn->conn_queue_entry = qe;
conn->conn_queue = conn->instance->connections;
queue_data_put_with_index(conn->instance->connections, qe);
}
conn->peer_node_id = peer_node_id;
etcp_update_log_name(conn);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] peer_node_id set to 0x%llx (state=%d reindexed=%d)",
conn->log_name, (unsigned long long)peer_node_id, conn->state, conn->state == 1);
}
// Update log_name when peer_node_id becomes known // Update log_name when peer_node_id becomes known

1
src/transport_layer/etcp.h

@ -263,7 +263,6 @@ struct ETCP_CONN {
// Functions // Functions
struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name); struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* name);
void etcp_connection_close(struct ETCP_CONN* etcp); void etcp_connection_close(struct ETCP_CONN* etcp);
void etcp_conn_set_peer_node_id(struct ETCP_CONN* conn, uint64_t peer_node_id);
void etcp_conn_queue_set_ready(struct ETCP_CONN* conn); // move pending->connections, fire ready cbks void etcp_conn_queue_set_ready(struct ETCP_CONN* conn); // move pending->connections, fire ready cbks
/** /**

166
src/transport_layer/etcp_connect.c

@ -43,14 +43,13 @@ static struct ETCP_CONNECT* connect_find(struct UTUN_INSTANCE* inst, uint64_t no
return NULL; return NULL;
} }
static void connect_deliver(struct ETCP_CONNECT* ctx, int type) { static void connect_deliver(struct ETCP_CONNECT* ctx, int type, int free_nodes) {
struct etcp_connect_cb_node** pp = &ctx->cb_list; struct etcp_connect_cb_node** pp = &ctx->cb_list;
while (*pp) { while (*pp) {
struct etcp_connect_cb_node* node = *pp; struct etcp_connect_cb_node* node = *pp;
if (type == 0 || (node->flags & (uint8_t)type)) { if (type == 0 || (node->flags & (uint8_t)type)) {
*pp = node->next;
node->cb(node->arg, (type == 0) ? NULL : ctx->conn, type); node->cb(node->arg, (type == 0) ? NULL : ctx->conn, type);
u_free(node); if (free_nodes) { *pp = node->next; u_free(node); } else { pp = &node->next; }
} else { } else {
pp = &node->next; pp = &node->next;
} }
@ -137,7 +136,7 @@ static void tcp_link_close_cb(struct stcp_link* link, int err, void* arg) {
static void connect_bgp_ready_cb(struct ETCP_CONN* conn) { static void connect_bgp_ready_cb(struct ETCP_CONN* conn) {
struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id); struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id);
if (!ctx || ctx->done) return; if (!ctx || ctx->done) return;
connect_deliver(ctx, ETCP_CONNECT_BGP_READY); connect_deliver(ctx, ETCP_CONNECT_BGP_READY, 1);
} }
static void connect_initial_timeout_cb(void* arg) { static void connect_initial_timeout_cb(void* arg) {
@ -148,7 +147,7 @@ static void connect_initial_timeout_cb(void* arg) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] timeout for node 0x%016llx", (unsigned long long)ctx->node_id); DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] timeout for node 0x%016llx", (unsigned long long)ctx->node_id);
struct ETCP_CONN* conn = ctx->conn; ctx->conn = NULL; struct ETCP_CONN* conn = ctx->conn; ctx->conn = NULL;
if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; } if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; }
connect_deliver(ctx, 0); connect_deliver(ctx, 0, 1);
if (conn) etcp_connection_close(conn); if (conn) etcp_connection_close(conn);
connect_cancel(ctx); connect_cancel(ctx);
} }
@ -176,24 +175,41 @@ static void connect_settle_timeout_cb(void* arg) {
stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL;
} }
etcp_conn_remove_init_cbk(ctx->conn, connect_init_cb, ctx); etcp_conn_remove_init_cbk(ctx->conn, connect_init_cb, ctx);
connect_deliver(ctx, ETCP_CONNECT_LATE); connect_deliver(ctx, ETCP_CONNECT_LATE, 1);
connect_cancel(ctx); connect_cancel(ctx);
} }
static void connect_init_cb(struct ETCP_CONN* conn, void* arg) { static void connect_init_cb(struct ETCP_CONN* conn, void* arg) {
struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg; struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg;
if (!ctx || ctx->done) return; if (!ctx || ctx->done) return;
int total_links = 0, up_links = 0;
struct ETCP_LINK* link = conn->links; struct ETCP_LINK* link = conn->links;
while (link) { while (link) {
total_links++;
if (link->initialized && link->rtt_last && link->rtt_last < ctx->min_rtt) if (link->initialized && link->rtt_last && link->rtt_last < ctx->min_rtt)
ctx->min_rtt = link->rtt_last; ctx->min_rtt = link->rtt_last;
if (link->link_state == 3) up_links++;
link = link->next; link = link->next;
} }
if (total_links == 1 || up_links == total_links || (total_links == 0 && conn->transport_link)) {
if (!ctx->early_delivered) {
ctx->early_delivered = 1;
if (ctx->initial_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->initial_timer); ctx->initial_timer = NULL; }
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] immediate LATE for node 0x%016llx (links=%d up=%d)",
(unsigned long long)ctx->node_id, total_links, up_links);
{ struct etcp_connect_cb_node* n = ctx->cb_list; while (n) { struct etcp_connect_cb_node* next = n->next; if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY); if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE); u_free(n); n = next; } ctx->cb_list = NULL; }
connect_cancel(ctx);
return;
}
if (ctx->early_delivered) return; if (ctx->early_delivered) return;
ctx->early_delivered = 1; ctx->early_delivered = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] early ready for node 0x%016llx, min_rtt=%u", DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] early ready for node 0x%016llx, min_rtt=%u (links=%d up=%d)",
(unsigned long long)ctx->node_id, ctx->min_rtt); (unsigned long long)ctx->node_id, ctx->min_rtt, total_links, up_links);
connect_deliver(ctx, ETCP_CONNECT_EARLY); connect_deliver(ctx, ETCP_CONNECT_EARLY, 0);
if (ctx->initial_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->initial_timer); ctx->initial_timer = NULL; } if (ctx->initial_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->initial_timer); ctx->initial_timer = NULL; }
uint32_t settle_tb; uint32_t settle_tb;
if (ctx->min_rtt == 0xFFFF) settle_tb = 20000; if (ctx->min_rtt == 0xFFFF) settle_tb = 20000;
@ -209,28 +225,35 @@ static void connect_init_cb(struct ETCP_CONN* conn, void* arg) {
int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node, int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags) { etcp_connect_callback_t cb, void* arg, uint8_t flags) {
if (!inst || !node || !cb) return -1; if (!inst || !node || !cb) return -1;
uint64_t node_id = node->node->node_id;
// Check if already connected uint64_t node_id = node->node->node_id;
{ if (node_id == 0) {
struct ETCP_CONN* existing = instance_find_conn(inst, node_id); if (!node->node->public_key) {
if (existing && existing->peer_node_id == node_id) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node_id=0 and public_key is NULL");
struct ETCP_LINK* l = existing->links; return -1;
while (l) { if (l->link_state == 3) break; l = l->next; }
if (l) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx already connected, delivering immediately",
(unsigned long long)node_id);
if (flags & ETCP_CONNECT_EARLY) cb(arg, existing, ETCP_CONNECT_EARLY);
if (flags & ETCP_CONNECT_LATE) cb(arg, existing, ETCP_CONNECT_LATE);
if ((flags & ETCP_CONNECT_BGP_READY) && existing->routing_exchange_active >= 3)
cb(arg, existing, ETCP_CONNECT_BGP_READY);
return 0;
} }
int zero = 1;
for (int i = 0; i < SC_PUBKEY_SIZE; i++)
if (node->node->public_key[i]) { zero = 0; break; }
if (zero) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node_id=0 and public_key is all zeros");
return -1;
} }
node_id = sc_derive_node_id_from_pubkey(node->node->public_key);
node->node->node_id = node_id;
} }
struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
struct ETCP_CONNECT* ctx = connect_find(inst, node_id); struct ETCP_CONNECT* ctx = connect_find(inst, node_id);
if (ctx) {
// Case 1: conn && ctx — already connecting
if (conn && ctx) {
if (ctx->done) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx: existing ctx timed out, cb(NULL,0)",
(unsigned long long)node_id);
cb(arg, NULL, 0);
return 0;
}
struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node)); struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) return -1; if (!cn) return -1;
cn->cb = cb; cn->arg = arg; cn->flags = flags; cn->next = ctx->cb_list; ctx->cb_list = cn; cn->cb = cb; cn->arg = arg; cn->flags = flags; cn->next = ctx->cb_list; ctx->cb_list = cn;
@ -239,9 +262,88 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node,
return 0; return 0;
} }
struct ETCP_CONN* conn = etcp_connection_create(inst, NULL); // Case 2: !conn && ctx — context exists but conn was removed (rare)
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] failed to create ETCP_CONN for node 0x%016llx", if (!conn && ctx) {
(unsigned long long)node_id); cb(arg, NULL, 0); return -1; } struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) return -1;
cn->cb = cb; cn->arg = arg; cn->flags = flags; cn->next = ctx->cb_list; ctx->cb_list = cn;
return 0;
}
// Case 3: conn && !ctx — existing connection, no context
if (conn && !ctx) {
if (conn->state == 0) {
// Pending conn from another path (e.g. incoming INIT before etcp_connect)
ctx = u_calloc(1, sizeof(struct ETCP_CONNECT));
if (!ctx) return -1;
ctx->instance = inst; ctx->node_id = node_id; ctx->conn = conn;
ctx->min_rtt = 0xFFFF;
struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) { u_free(ctx); return -1; }
cn->cb = cb; cn->arg = arg; cn->flags = flags; ctx->cb_list = cn;
etcp_conn_add_init_cbk(conn, connect_init_cb, ctx);
conn->bgp_ready_cbk = connect_bgp_ready_cb;
ctx->initial_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, ctx,
connect_initial_timeout_cb, "etcp_connect_init");
ctx->next = inst->pending_connects; inst->pending_connects = ctx;
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx: existing pending conn, waiting for init",
(unsigned long long)node_id);
return 0;
}
// conn->state == 1: already ready — deliver EARLY+LATE immediately (via direct loop to deliver both phases)
ctx = u_calloc(1, sizeof(struct ETCP_CONNECT));
if (!ctx) return -1;
ctx->instance = inst; ctx->node_id = node_id; ctx->conn = conn;
ctx->early_delivered = 1;
{
struct etcp_connect_cb_node* cn = u_calloc(1, sizeof(struct etcp_connect_cb_node));
if (!cn) { u_free(ctx); return -1; }
cn->cb = cb; cn->arg = arg; cn->flags = flags;
ctx->cb_list = cn;
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx: existing ready conn, delivering EARLY+LATE",
(unsigned long long)node_id);
struct etcp_connect_cb_node* n = ctx->cb_list;
struct etcp_connect_cb_node* keep = NULL;
while (n) {
struct etcp_connect_cb_node* next = n->next;
if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY);
if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE);
if (n->flags & ETCP_CONNECT_BGP_READY) {
if (conn->routing_exchange_active >= 3) n->cb(n->arg, ctx->conn, ETCP_CONNECT_BGP_READY);
else { n->next = keep; keep = n; n = next; continue; }
}
u_free(n);
n = next;
}
ctx->cb_list = keep;
if (keep) {
conn->bgp_ready_cbk = connect_bgp_ready_cb;
ctx->next = inst->pending_connects; inst->pending_connects = ctx;
} else {
u_free(ctx);
}
return 0;
}
// Case 4: !conn && !ctx — new connection
conn = etcp_connection_create(inst, NULL);
if (!conn) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] failed to create ETCP_CONN for node 0x%016llx",
(unsigned long long)node_id);
cb(arg, NULL, 0);
return -1;
}
conn->peer_node_id = node_id;
etcp_update_log_name(conn);
if (conn->conn_queue && conn->conn_queue_entry) {
queue_remove_data(conn->conn_queue, conn->conn_queue_entry);
queue_entry_free(conn->conn_queue_entry);
}
{ struct ll_entry* qe = queue_entry_new(sizeof(struct conn_queue_entry)); if (qe) { struct conn_queue_entry* ce = (struct conn_queue_entry*)qe->data; ce->peer_node_id = node_id; ce->conn = conn; conn->conn_queue_entry = qe; conn->conn_queue = inst->connections; queue_data_put_with_index(inst->connections, qe); } }
if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) { if (sc_init_ctx(&conn->crypto_ctx, &inst->my_keys) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] sc_init_ctx failed for node 0x%016llx", (unsigned long long)node_id); DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] sc_init_ctx failed for node 0x%016llx", (unsigned long long)node_id);
etcp_connection_close(conn); cb(arg, NULL, 0); return -1; etcp_connection_close(conn); cb(arg, NULL, 0); return -1;
@ -252,8 +354,10 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node,
} }
ctx = u_calloc(1, sizeof(struct ETCP_CONNECT)); ctx = u_calloc(1, sizeof(struct ETCP_CONNECT));
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] alloc failed"); if (!ctx) {
etcp_connection_close(conn); cb(arg, NULL, 0); return -1; } DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] alloc failed");
etcp_connection_close(conn); cb(arg, NULL, 0); return -1;
}
ctx->instance = inst; ctx->node_id = node_id; ctx->conn = conn; ctx->instance = inst; ctx->node_id = node_id; ctx->conn = conn;
ctx->min_rtt = 0xFFFF; ctx->min_rtt = 0xFFFF;
{ {
@ -272,7 +376,7 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_NODEQ* node,
DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] no links created for node 0x%016llx", DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] no links created for node 0x%016llx",
(unsigned long long)node_id); (unsigned long long)node_id);
if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; } if (ctx->tcp_link) { stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; }
connect_deliver(ctx, 0); connect_deliver(ctx, 0, 1);
etcp_connection_close(conn); u_free(ctx); return -1; etcp_connection_close(conn); u_free(ctx); return -1;
} }

156
src/transport_layer/etcp_connect.h

@ -1,23 +1,146 @@
#ifndef ETCP_CONNECT_H #ifndef ETCP_CONNECT_H
#define ETCP_CONNECT_H #define ETCP_CONNECT_H
#include <stdint.h>
struct UTUN_INSTANCE; struct UTUN_INSTANCE;
struct ETCP_CONN; struct ETCP_CONN;
struct TOPO_NODEQ;
typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int type);
/* /**
* ============================================================================ * ============================================================================
* Штатное закрытие исходящих соединений (etcp_connect) * etcp_connect() — асинхронное подключение к удалённому узлу
* ============================================================================ * ============================================================================
* *
* Исходящее соединение создаётся через etcp_connect(): * Инициирует (или находит существующее) ETCP-соединение к узлу.
* etcp_connect(inst, node, cb, arg, flags) * Вызов неблокирующий: коллбэк cb вызывается асинхронно по мере прохождения
* ├── etcp_connection_create(inst, NULL) → conn (добавляется в inst->connections) * фаз установки соединения.
* ├── ctx = u_calloc(...); ctx->conn = conn *
* ├── для TCP: ctx->tcp_link = stcp_link_connect(...) * @param inst UTUN_INSTANCE
* │ conn->transport_link = ctx->tcp_link * @param node узел-адресат (TOPO_NODEQ). Из node->node берутся:
* └── ctx->next = inst->pending_connects; inst->pending_connects = ctx * - node_id — идентификатор узла (0 = вычислить из pubkey)
* - public_key — X25519 pubkey (32 байта)
* - v4_addrs / v6_addrs — адреса для создания линков
* @param cb коллбэк (etcp_connect_callback_t)
* @param arg пользовательский аргумент, передаваемый в cb
* @param flags битовая маска интересующих фаз: ETCP_CONNECT_EARLY(1) |
* ETCP_CONNECT_LATE(2) | ETCP_CONNECT_BGP_READY(4)
* @return 0 при успехе, -1 при ошибке (до вызова коллбэка)
*
*
* ---- Деривация node_id ----
*
* Если node->node->node_id == 0, он вычисляется из public_key:
* node_id = SHA256(pubkey, 32)[0..7] & 0x7FFFFFFFFFFFFFFF
*
* Требования к ключу:
* - public_key не NULL
* - не все 32 байта нулевые
* Иначе — возврат -1 без вызова коллбэка.
*
* Вычисленный node_id сохраняется обратно в node->node->node_id.
*
*
* ---- Индексация соединения ----
* *
* Пути завершения: * Сразу после etcp_connection_create() conn->peer_node_id устанавливается
* в node_id, что позволяет etcp_conn_queue_set_ready() (при инициализации
* первого линка) переиндексировать запись в очереди inst->connections
* с key=0 на key=node_id.
*
* До переиндексации соединение находится в очереди с key=0 и доступно
* только через connect_find() в pending_connects или по совпадению адреса.
*
*
* ---- Поиск существующего соединения ----
*
* Перед созданием нового conn выполняется поиск:
*
* a. instance_find_conn(inst, node_id) — активное соединение в очереди
* inst->connections (UDP) или inst->tcp_connections (TCP)
* b. connect_find(inst, node_id) — контекст в inst->pending_connects
*
* Возможные комбинации:
*
* 1. conn && ctx — уже подключается
* ├─ ctx->done (таймаут уже сработал) → cb(arg, NULL, 0) сразу
* └─ иначе → cb добавляется в ctx->cb_list (ждёт своей фазы)
*
* 2. !conn && ctx — контекст есть, но conn был удалён (редкий случай)
* → cb добавляется в ctx->cb_list
*
* 3. conn && !ctx — соединение рабочее, контекста нет (входящее/серверное)
* ├─ conn->state == 0 (pending, ещё не инициализирован):
* │ → создаётся ctx с connect_init_cb, initial_timer
* │ → ctx добавляется в pending_connects
* │ → ожидание инициализации в обычном порядке
* │
* └─ conn->state == 1 (ready, уже проинициализирован):
* → connect_deliver(ctx, ETCP_CONNECT_LATE) сразу
* → если routing_exchange_active >= 3:
* connect_deliver(ctx, ETCP_CONNECT_BGP_READY)
* → если в cb_list остались неудовлетворённые коллбэки
* (ждут BGP_READY, который ещё не готов):
* ctx остаётся в pending_connects с conn->bgp_ready_cbk
*
* 4. !conn && !ctx — новое подключение (основной путь)
* → etcp_connection_create(inst, NULL)
* → conn->peer_node_id = node_id (индексация deferred до queue_set_ready)
* → sc_init_ctx + sc_set_peer_public_key — настройка крипто
* → connect_create_links_v4/v6 — создание UDP + TCP линков
* → initial_timer = uasync_set_timeout(connect_timeout_tb, ...)
* → ctx добавляется в pending_connects
*
*
* ---- Фазы доставки коллбэков ----
*
* connect_init_cb (вызывается при conn->initialized, первый линк отдал рукопожатие):
*
* Подсчитывается total_links и up_links (link_state == 3):
*
* ├─ total_links == 1: один линк — EARLY не нужен
* │ → connect_deliver(ctx, ETCP_CONNECT_LATE) — сразу LATE
* │ → ctx освобождается
* │
* ├─ up_links == total_links: все линки уже UP (link_state==3)
* │ → connect_deliver(ctx, ETCP_CONNECT_LATE) — сразу LATE
* │ → ctx освобождается
* │
* └─ иначе (часть линков UP, часть ещё в handshake):
* → connect_deliver(ctx, ETCP_CONNECT_EARLY) — уведомить о доступности
* → initial_timer отменяется
* → запускается settle_timer = min_rtt * 8 (clamped [500ms, 2s] в 0.1ms)
*
* connect_settle_timeout_cb (settle-таймер):
* ├─ Закрывает неинициализированные линки (!link->initialized)
* ├─ Закрывает TCP-линк если !tcp_ready
* ├─ Убирает connect_init_cb из conn->init_cbks
* └─ connect_deliver(ctx, ETCP_CONNECT_LATE) → ctx освобождается
*
* connect_bgp_ready_cb (вызывается при routing_exchange_active >= 3):
* └─ connect_deliver(ctx, ETCP_CONNECT_BGP_READY)
* (может произойти в любой момент после установки соединения)
*
* connect_initial_timeout_cb (initial_timer):
* ├─ ctx->done = 1
* ├─ Закрывает TCP-линк
* ├─ connect_deliver(ctx, 0) → cb(arg, NULL, 0) для всех
* ├─ etcp_connection_close(conn) — закрывает соединение
* └─ ctx освобождается
*
*
* ---- Ошибки ----
*
* Коллбэк вызывается с type=0, conn=NULL в случаях:
* - initial_timer истёк (ни один линк не инициализировался)
* - все линки были закрыты до завершения handshake
* - ошибка аллокации / крипто на этапе создания
*
*
* ---- Штатное закрытие исходящих соединений ----
*
* Пути завершения ctx:
* *
* 1. Таймаут установки (connect_initial_timeout_cb): * 1. Таймаут установки (connect_initial_timeout_cb):
* stcp_link_close(tcp_link) → etcp_connection_close(conn) → connect_cancel(ctx) * stcp_link_close(tcp_link) → etcp_connection_close(conn) → connect_cancel(ctx)
@ -31,21 +154,18 @@ struct ETCP_CONN;
* │ Фаза 1 (detach): закрыть линки, удалить из inst->connections, state=2 * │ Фаза 1 (detach): закрыть линки, удалить из inst->connections, state=2
* │ Фаза 2 (deferred): uasync_call_soon(ua, conn, etcp_connection_free_deferred) * │ Фаза 2 (deferred): uasync_call_soon(ua, conn, etcp_connection_free_deferred)
* │ * │
* ├── while (immediate_queue_head) uasync_poll(ua, 0) // ждём завершения * ├── while (immediate_queue_head) uasync_poll(ua, 0)
* │ └── etcp_connection_free_resources(conn) * │ └── etcp_connection_free_resources(conn)
* │ ├── drain_and_free_queue(*) * │ ├── drain_and_free_queue(*)
* │ ├── etcp_connect_cancel_for_conn(inst, conn) ← эта функция * │ ├── etcp_connect_cancel_for_conn(inst, conn)
* │ └── u_free(etcp) * │ └── u_free(etcp)
* │ * │
* └── memory_pool_destroy(pkt/ack/data) // теперь безопасно * └── memory_pool_destroy(pkt/ack/data)
*
* Закрытие tcp_link НЕ входит в ответственность cancel-функций.
* tcp_link закрывается ДО вызова connect_cancel / cancel_for_conn
* (в путях 1 и 2 явно; в пути 3 tcp_link закрыт ранее или утекает — pre-existing).
* *
* ============================================================================ * ============================================================================
*/ */
int etcp_connect(struct UTUN_INSTANCE* instance, struct TOPO_NODEQ* node,
etcp_connect_callback_t cb, void* arg, uint8_t flags);
/* /*
* ============================================================================ * ============================================================================

6
src/transport_layer/etcp_connections.c

@ -1523,7 +1523,8 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D
DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "[%s] Received Ed25519 pubkey from peer", link->etcp->log_name); DEBUG_DEBUG(DEBUG_CATEGORY_CRYPTO, "[%s] Received Ed25519 pubkey from peer", link->etcp->log_name);
etcp_conn_set_peer_node_id(link->etcp, server_node_id); link->etcp->peer_node_id = server_node_id;
etcp_update_log_name(link->etcp);
link->initialized = 1;// получен init response (client) link->initialized = 1;// получен init response (client)
link->link_state = 3; // connected link->link_state = 3; // connected
@ -1798,7 +1799,8 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) {
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[1] : 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[2] : 0,
conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[3] : 0); conn->crypto_ctx.pk ? conn->crypto_ctx.pk->public_key[3] : 0);
etcp_conn_set_peer_node_id(conn, peer_id); 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); 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);
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p total=%d", DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p total=%d",
ip_to_str(&addr, addr.ss_family).str, peer_id, conn, ip_to_str(&addr, addr.ss_family).str, peer_id, conn,

8
tests/test_uasync_socket_race.c

@ -60,16 +60,18 @@ static void on_closed_cb(struct tcp_conn* tc, void* arg) {
static void on_accept_cb(int fd, void* arg) { static void on_accept_cb(int fd, void* arg) {
(void)arg; (void)arg;
while (1) {
struct sockaddr_in addr; socklen_t alen = sizeof(addr); struct sockaddr_in addr; socklen_t alen = sizeof(addr);
int csock = accept(fd, (struct sockaddr*)&addr, &alen); int csock = accept(fd, (struct sockaddr*)&addr, &alen);
if (csock < 0) return; if (csock < 0) return;
if (csock == 0) { int r = open("/dev/null", O_RDONLY); if (r > 0 && r != 0) { dup2(r, 0); close(r); } return; } if (csock == 0) { int r = open("/dev/null", O_RDONLY); if (r > 0 && r != 0) { dup2(r, 0); close(r); } continue; }
socket_set_nonblocking(csock); socket_set_nonblocking(csock);
struct tcp_conn* tc = tcp_conn_create(g_ua, csock, 512, 512, 4, 0, 0, on_fin_cb, on_error_cb, NULL); struct tcp_conn* tc = tcp_conn_create(g_ua, csock, 512, 512, 4, 0, 0, on_fin_cb, on_error_cb, NULL);
if (!tc) { socket_close_wrapper(csock); return; } if (!tc) { socket_close_wrapper(csock); continue; }
tc->on_closed = on_closed_cb; tc->on_closed = on_closed_cb;
queue_set_callback(tc->read_queue, on_read_cb, tc); queue_set_callback(tc->read_queue, on_read_cb, tc);
g_conn_count++; g_conn_count++;
}
} }
static int child_main(int port, int id, int count) { static int child_main(int port, int id, int count) {
@ -97,7 +99,7 @@ static void monitor(void* arg) {
if (WIFEXITED(status) && WEXITSTATUS(status) != 0) { g_done = -1; return; } if (WIFEXITED(status) && WEXITSTATUS(status) != 0) { g_done = -1; return; }
g_children[i] = 0; g_children[i] = 0;
} }
if (alive == 0) { g_ok = 1; g_done = 1; } if (alive == 0 && g_read_count >= CHILDREN * ITER_PER_CHILD) { g_ok = 1; g_done = 1; }
if (!g_done) g_mon_id = uasync_set_timeout(g_ua, 500, NULL, monitor, "mon"); if (!g_done) g_mon_id = uasync_set_timeout(g_ua, 500, NULL, monitor, "mon");
} }

Loading…
Cancel
Save