diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 3214fc7e..9fd4bec5 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -1578,7 +1578,7 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { errorcode=2; goto ec_fr; } - DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "INIT X25519 OK from %s", sockaddr_storage_to_str(&addr).str); + DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "X25519 decrypt OK from %s", sockaddr_storage_to_str(&addr).str); if (sc_decrypt(&sc, data, recv_len - SC_PUBKEY_ENC_SIZE, (uint8_t*)&pkt->timestamp, &pkt_len)) { DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "failed to decrypt init packet, from %s", sockaddr_storage_to_str(&addr).str); errorcode=3; @@ -1602,14 +1602,14 @@ 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) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "INIT decrypted: PING from peer=0x%016llx src=%s", + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "X25519 decrypted: PING from peer=0x%016llx src=%s", (unsigned long long)peer_id, sockaddr_storage_to_str(&addr).str); int ret = handle_ping(e_sock, pkt, &addr, decrypted_pubkey, pkt_len); if (ret) { errorcode = ret; goto ec_fr; } return; } if (code == ETCP_PONG) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "INIT decrypted: PONG from peer=0x%016llx src=%s", + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "X25519 decrypted: PONG from peer=0x%016llx src=%s", (unsigned long long)peer_id, sockaddr_storage_to_str(&addr).str); int ret = handle_pong(e_sock, pkt, &addr, pkt_len); if (ret) { errorcode = ret; goto ec_fr; } @@ -1727,9 +1727,9 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { // For CHANNEL_INIT (0x04): if link already initialized - no reset, otherwise reset // For INIT_REQUEST (0x02): always reset // Check session_id: if same - no reinit, if different - client restarted, do reinit - if (code == ETCP_INIT_REQUEST || conn->session_id != session_id) { - send_reset = 1; // Client explicitly requested reset, or new session - DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit (code=%02x session was %08x now %08x)", code, conn->session_id, session_id); + if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || new_conn) { + send_reset = 1; // Client explicitly requested reset, or new session, or new server-side connection + DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit (code=%02x session was %08x now %08x new_conn=%d)", code, conn->session_id, session_id, new_conn); conn->session_id = session_id; etcp_conn_reinit(conn); } else { @@ -1753,9 +1753,9 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { link->remote_type = req->type; // For new links: reset if client requested or session changed - if (code == ETCP_INIT_REQUEST || conn->session_id != session_id) { + if (code == ETCP_INIT_REQUEST || conn->session_id != session_id || new_conn) { send_reset = 1; - DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit for new link (code=%02x session was %08x now %08x)", code, conn->session_id, session_id); + DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "do reinit for new link (code=%02x session was %08x now %08x new_conn=%d)", code, conn->session_id, session_id, new_conn); conn->session_id = session_id; etcp_conn_reinit(conn); } diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index eaa848f6..5edd7ae4 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -16,6 +16,7 @@ #include "../../../src/topo_group.h" #include "../../../src/topo_node.h" #include "../../../src/conn_mgr.h" +#include "../../../src/etcp.h" #include "../../../src/etcp_connections.h" #include "../../../src/secure_channel.h" #include "../../../lib/u_async.h" @@ -584,6 +585,80 @@ static void connect_result_cb(int result, uint64_t node_id, void* arg) { gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); } +#define CC_PARALLEL_CONNECT_TIMEOUT_MS 3000 + +struct cc_parallel_state { + struct UTUN_INSTANCE* inst; + int addr_count; + int pending_count; + int completed; /* 0=pending, 1=success, -1=fail-delivered */ + uint64_t node_id; + uint64_t channel_id; + struct ETCP_CONN** conns; + void** timers; +}; + +struct cc_parallel_ctx { + struct cc_parallel_state* state; + int addr_index; +}; + +static void cc_parallel_cleanup(struct cc_parallel_state* st) { + if (!st) return; + u_free(st->conns); + u_free(st->timers); + u_free(st); +} + +static void cc_parallel_ready_cb(struct ETCP_CONN* conn, void* arg) { + struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg; + struct cc_parallel_state* st = pctx->state; + if (st->completed) { u_free(pctx); return; } + st->completed = 1; + + uint8_t data[20]; + memcpy(data, &st->node_id, 8); + int r = CONN_MGR_OK; + memcpy(data + 8, &r, 4); + memcpy(data + 12, &st->channel_id, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: parallel connect SUCCESS idx=%d peer=0x%016llx", + CC_ID, pctx->addr_index, (unsigned long long)st->node_id); + + for (int i = 0; i < st->addr_count; i++) { + if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; } + if (st->conns[i] && i != pctx->addr_index) { etcp_connection_close(st->conns[i]); st->conns[i] = NULL; } + } + u_free(pctx); + cc_parallel_cleanup(st); +} + +static void cc_parallel_timeout_cb(void* arg) { + struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg; + struct cc_parallel_state* st = pctx->state; + if (st->completed) { u_free(pctx); return; } + + if (st->conns[pctx->addr_index]) { etcp_connection_close(st->conns[pctx->addr_index]); st->conns[pctx->addr_index] = NULL; } + st->timers[pctx->addr_index] = NULL; + st->pending_count--; + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: parallel connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx", + CC_ID, pctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id); + + if (st->pending_count <= 0 && !st->completed) { + st->completed = -1; + uint8_t data[20]; + memcpy(data, &st->node_id, 8); + int r = CONN_MGR_ERR_TIMEOUT; + memcpy(data + 8, &r, 4); + memcpy(data + 12, &st->channel_id, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); + cc_parallel_cleanup(st); + } + u_free(pctx); +} + void chat_core_connect_from_invite(struct chat_invite* inv) { if (!g_cc.initialized || !g_cc.inst || !inv) return; struct TOPO_GROUP* group = topo_groups_get_default(g_cc.inst->topo_groups); @@ -655,10 +730,89 @@ void chat_core_connect_from_invite(struct chat_invite* inv) { queue_data_put_with_index(group->nodes, &nq->ll); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: created NODEINFO for 0x%016llx, %d addrs, connecting", + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: created NODEINFO for 0x%016llx, %d addrs, parallel connect", CC_ID, (unsigned long long)node_id, inv->addr_count); - conn_mgr_connect_node(g_cc.inst->conn_mgr, node_id, 30000, - connect_result_cb, &inv->channel_id); + + /* count IPv4 addresses and collect them for parallel connect */ + int v4_count = 0; + for (struct TOPO_ADDR4* a = addrs_head; a; a = a->next) v4_count++; + if (v4_count == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: no IPv4 addresses in invite", CC_ID); + uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_NO_ADDRESSES; + memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); + return; + } + + struct ETCP_SOCKET* best_socket = g_cc.inst->etcp_sockets; + while (best_socket && best_socket->local_addr.ss_family != AF_INET) best_socket = best_socket->next; + if (!best_socket) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: no AF_INET socket", CC_ID); + uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL; + memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); + return; + } + + struct cc_parallel_state* pst = u_calloc(1, sizeof(struct cc_parallel_state)); + if (!pst) { + uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL; + memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); + return; + } + pst->inst = g_cc.inst; + pst->addr_count = v4_count; + pst->pending_count = v4_count; + pst->completed = 0; + pst->node_id = node_id; + pst->channel_id = inv->channel_id; + pst->conns = u_calloc(v4_count, sizeof(struct ETCP_CONN*)); + pst->timers = u_calloc(v4_count, sizeof(void*)); + if (!pst->conns || !pst->timers) { cc_parallel_cleanup(pst); return; } + + int idx = 0; + for (struct TOPO_ADDR4* a = addrs_head; a && idx < v4_count; a = a->next, idx++) { + 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); + struct sockaddr_storage sa; + memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); + + struct ETCP_CONN* conn = etcp_connection_create(g_cc.inst, NULL); + if (!conn) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: etcp_connection_create failed idx=%d", CC_ID, idx); + pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; + } + sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); + sc_set_peer_public_key(&conn->crypto_ctx, inv->pubkey, 0); + + struct cc_parallel_ctx* pctx = u_calloc(1, sizeof(struct cc_parallel_ctx)); + if (!pctx) { etcp_connection_close(conn); pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; } + pctx->state = pst; pctx->addr_index = idx; + + etcp_conn_set_ready_cbk(conn, cc_parallel_ready_cb, pctx); + + if (!etcp_link_new(conn, best_socket, &sa, 0)) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: etcp_link_new failed idx=%d", CC_ID, idx); + u_free(pctx); etcp_connection_close(conn); + pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; + } + pst->conns[idx] = conn; + pst->timers[idx] = uasync_set_timeout(g_cc.inst->ua, CC_PARALLEL_CONNECT_TIMEOUT_MS * 10, + pctx, cc_parallel_timeout_cb, "cc_parallel"); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: parallel connect attempt %d/%d to %d.%d.%d.%d:%d", + CC_ID, idx + 1, v4_count, sin.sin_addr.s_addr & 0xFF, (sin.sin_addr.s_addr >> 8) & 0xFF, + (sin.sin_addr.s_addr >> 16) & 0xFF, (sin.sin_addr.s_addr >> 24) & 0xFF, a->port); + } + + if (pst->pending_count <= 0 && !pst->completed) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "%s: all %d parallel connects failed to start", CC_ID, v4_count); + uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_UNREACHABLE; + memcpy(err + 8, &r, 4); memcpy(err + 12, &inv->channel_id, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); + cc_parallel_cleanup(pst); + } } /* ─── создание канала ─── */