diff --git a/src/routing_layer/topo_group_connect.c b/src/routing_layer/topo_group_connect.c index 12376bf4..ddb957d3 100644 --- a/src/routing_layer/topo_group_connect.c +++ b/src/routing_layer/topo_group_connect.c @@ -45,6 +45,7 @@ struct TOPO_GROUP_CONNECT { int connected_count; int active_conn_count; int cursor; + int tried_super; uint64_t* candidate_ids; int candidate_count; struct CONN_MGR_HANDLE* handles[TGC_MAX_HANDLES]; @@ -67,7 +68,7 @@ int topo_group_connect_init(struct TOPO_GROUP* group) { return -1; struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc)); if (!gc) return -1; - gc->group = group; gc->active = 1; + gc->group = group; gc->active = 1; gc->tried_super = 0; group->connect = gc; uint64_t* ids = NULL; int count = 0; @@ -277,24 +278,37 @@ static void tgc_phase1_timeout(void* arg) { static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) { if (gc->phase != TGC_PHASE_TWO) return; - if (gc->candidate_count == 0) { - sqlite3* db = gc->group->instance->topo_sqlite_db; - topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count, - gc->group->instance->node_id); - gc->cursor = 0; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 loaded %d public peers for ch=%s", - TGC_ID, gc->candidate_count, gc->group->channel_id); - } - while (gc->cursor < gc->candidate_count) { - uint64_t nid = gc->candidate_ids[gc->cursor++]; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d)", TGC_ID, - (unsigned long long)nid, gc->cursor, gc->candidate_count); - conn_mgr_open_invite(gc->group->instance, gc->group->group_id, NULL, nid, tgc_callback, gc, NULL); - return; + while (1) { + if (gc->candidate_count == 0) { + sqlite3* db = gc->group->instance->topo_sqlite_db; + if (gc->tried_super == 0) { + topo_node_sqlite_get_supernode_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count); + gc->tried_super = 1; gc->cursor = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 supernode round — loaded %d for ch=%s", + TGC_ID, gc->candidate_count, gc->group->channel_id); + } else if (gc->tried_super == 1) { + topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count, + gc->group->instance->node_id); + gc->tried_super = 2; gc->cursor = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 public round — loaded %d for ch=%s", + TGC_ID, gc->candidate_count, gc->group->channel_id); + } else { + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id); + gc->phase = TGC_PHASE_THREE; gc->candidate_count = 0; + tgc_phase3_try_next(gc); + return; + } + } + if (gc->cursor < gc->candidate_count) { + uint64_t nid = gc->candidate_ids[gc->cursor++]; + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d) super_round=%d", TGC_ID, + (unsigned long long)nid, gc->cursor, gc->candidate_count, gc->tried_super); + conn_mgr_open_invite(gc->group->instance, gc->group->group_id, NULL, nid, tgc_callback, gc, NULL); + return; + } + if (gc->candidate_ids) { u_free(gc->candidate_ids); gc->candidate_ids = NULL; } + gc->candidate_count = 0; } - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id); - gc->phase = TGC_PHASE_THREE; gc->candidate_count = 0; - tgc_phase3_try_next(gc); } /* ═══════════════════════════════════════════════════════════════════════ diff --git a/src/routing_layer/topo_node_sqlite.c b/src/routing_layer/topo_node_sqlite.c index 4c40de64..1e8cc539 100644 --- a/src/routing_layer/topo_node_sqlite.c +++ b/src/routing_layer/topo_node_sqlite.c @@ -747,6 +747,29 @@ int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, return 0; } +int topo_node_sqlite_get_supernode_peers(sqlite3* db, const char* channel_id, + uint64_t** out_ids, int* out_count) { + if (!db || !channel_id || !out_ids || !out_count) return -1; + *out_ids = NULL; *out_count = 0; + char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl)); + char sql[256]; snprintf(sql, sizeof(sql), + "SELECT p.node_id FROM \"%s\" p WHERE p.node_type=4" + " AND EXISTS (SELECT 1 FROM node_addresses na WHERE na.node_id=p.node_id AND na.family=4)", + peers_tbl); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1; + int cnt = 0; + while (sqlite3_step(st) == SQLITE_ROW) cnt++; + if (cnt == 0) { sqlite3_finalize(st); return 0; } + uint64_t* ids = u_malloc((size_t)cnt * sizeof(uint64_t)); + if (!ids) { sqlite3_finalize(st); return -1; } + sqlite3_reset(st); int i = 0; + while (sqlite3_step(st) == SQLITE_ROW) ids[i++] = (uint64_t)sqlite3_column_int64(st, 0); + sqlite3_finalize(st); + *out_ids = ids; *out_count = cnt; + return 0; +} + int topo_node_sqlite_get_local_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count, uint64_t self_node_id) { if (!db || !out_ids || !out_count) return -1; diff --git a/src/routing_layer/topo_node_sqlite.h b/src/routing_layer/topo_node_sqlite.h index 3464437e..a44dc8cb 100644 --- a/src/routing_layer/topo_node_sqlite.h +++ b/src/routing_layer/topo_node_sqlite.h @@ -70,5 +70,6 @@ int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count); int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count, uint64_t self_node_id); int topo_node_sqlite_get_local_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count, uint64_t self_node_id); +int topo_node_sqlite_get_supernode_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count); #endif diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index 7a47eef2..09c77ed6 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -309,8 +309,17 @@ void etcp_cbk_fire(struct ETCP_CONN* conn, int event) { static void etcp_on_up(struct ETCP_CONN* etcp) { int total_links = 0; for (struct ETCP_LINK* l = etcp->links; l; l = l->next) total_links++; uint64_t elapsed_ms = (get_time_tb() - etcp->setup_start_tb) / 10; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection UP (%d/%d links up, mtu=%d, setup=%llums, reinit=%u)", - etcp->log_name, etcp->links_up, total_links, etcp->mtu, (unsigned long long)elapsed_ms, etcp->reinit_count); + char links_str[256] = {0}; int pp = 0; + for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { + if (l->link_status) { + if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); + pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); + } + } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection UP (%d/%d links, mtu=%d, reinit=%u)%s%s", + etcp->log_name, etcp->links_up, total_links, etcp->mtu, etcp->reinit_count, + pp > 0 ? " — up: " : "", links_str); + (void)elapsed_ms; DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "CRYPTO_ON_UP_BEFORE: log=%s seskey=%02x%02x%02x%02x peer_pub=%02x%02x%02x%02x", etcp->log_name, etcp->crypto_ctx.session_key[0], etcp->crypto_ctx.session_key[1], @@ -328,7 +337,12 @@ static void etcp_on_up(struct ETCP_CONN* etcp) { } static void etcp_on_down(struct ETCP_CONN* etcp) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN (links_up=%d)", etcp->log_name, etcp->links_up); + char links_str[256] = {0}; int pp = 0; + for (struct ETCP_LINK* l = etcp->links; l; l = l->next) { + if (pp) pp += snprintf(links_str + pp, sizeof(links_str) - pp, ", "); + pp += snprintf(links_str + pp, sizeof(links_str) - pp, "%s", sockaddr_storage_to_str(&l->remote_addr).str); + } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Connection DOWN — links: %s", etcp->log_name, pp ? links_str : "(none)"); etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN); etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DOWN); } diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 0c321bb8..4f3286f5 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -336,7 +336,7 @@ static void keepalive_timer_cb(void* arg) { etcp_fire_link_status_cbk(link, link->link_state, old_link_status); etcp_on_link_down(link->etcp); if (old_link_status) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link down: link_id=%d state=%d init=%d ka=%d remote_ka=%d tmo: %llu>%llu els=%llums", link->etcp->log_name, link->local_link_id, link->link_state, link->initialized, link->recv_keepalive, link->remote_keepalive, (unsigned long long)timeout_units, (unsigned long long)elapsed, (unsigned long long)(elapsed/10)); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d down: addr=%s ka=%d remote_ka=%d state=%d init=%d tmo: %llu>%llu els=%llums", link->etcp->log_name, link->local_link_id, sockaddr_storage_to_str(&link->remote_addr).str, link->recv_keepalive, link->remote_keepalive, link->link_state, link->initialized, (unsigned long long)timeout_units, (unsigned long long)elapsed, (unsigned long long)(elapsed/10)); } } } @@ -1008,8 +1008,8 @@ void etcp_link_close(struct ETCP_LINK* link) { remove_link_from_queue(link); etcp_conn_on_inflight_lim_changed(link->etcp); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "[%s] Link %d closed: rcvd=%zub ack=%llub infl=%ub/%upkt", - link->etcp->log_name, link->local_link_id, link->total_decrypted, + 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); u_free(link->bbr); u_free(link); @@ -1496,9 +1496,9 @@ static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pk } if (link->etcp->initialized == 0) { etcp_conn_ready(link->etcp); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection established (socket=%s, link=%d)", link->etcp->log_name, e_sock->name, link->local_link_id); + 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)", link->etcp->log_name, link->local_link_id, link->mtu_local); + 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) @@ -1610,7 +1610,7 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D } loadbalancer_link_ready(link); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Link %d UP (client, mtu=%d)", link->etcp->log_name, link->local_link_id, link->mtu); + 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);