From 6edd40d1418c38af79f6cc9ef72451277308aed8 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Wed, 22 Jul 2026 22:48:02 +0300 Subject: [PATCH] =?UTF-8?q?topo=5Fgroup=5Fconnect:=20=D0=B0=D0=B2=D1=82?= =?UTF-8?q?=D0=BE-=D0=BF=D0=BE=D0=B4=D0=BA=D0=BB=D1=8E=D1=87=D0=B5=D0=BD?= =?UTF-8?q?=D0=B8=D0=B5=20CHAT-=D0=B3=D1=80=D1=83=D0=BF=D0=BF=20+=20conn?= =?UTF-8?q?=5Fmgr=5Fconnect=5Ffrom=5Finvite=20+=20=D0=B7=D0=B0=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D0=B0=20chat=5Fconn=5Fmgr=20=D0=B2=20chatgui?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Добавлено: - conn_mgr_connect_from_invite() — подключение к новому узлу по invite-данным (INIT → INVITE_INFO_REQ/RESP → проверка членства в группе → сохранение в БД) - topo_group_connect.c — авто-подключение при старте CHAT-группы: Phase 1: одновременный connect всех connected=1 пиров, таймаут, отмена незавершённых Phase 2: последовательный перебор публичных адресов до первого успеха - Счётчик active_conn_count с проверкой nodeinfo-путей (UP/DOWN) - restart() при phase=DONE и count==0 БД: - connected INTEGER NOT NULL DEFAULT 0 в peers_ + ALTER TABLE для существующих - topo_node_sqlite_set_connected / get_connected_peers / get_public_peers Chatgui: - Включён topo_group (убран utun_instance_set_topo_group_enabled(0)) - Создание TOPO_GROUP_TYPE_CHAT при channel_ready и при старте - Invite connect через conn_mgr_connect_from_invite - Удалён chat_conn_mgr (из сборки) + весь старый auto-connect код --- src/Makefile.am | 2 + src/routing_layer/conn_mgr.c | 289 +++++++++++++++++++++ src/routing_layer/conn_mgr.h | 55 ++++ src/routing_layer/topo_group.c | 6 + src/routing_layer/topo_group.h | 8 + src/routing_layer/topo_group_connect.c | 260 +++++++++++++++++++ src/routing_layer/topo_group_connect.h | 30 +++ src/routing_layer/topo_node_sqlite.c | 63 ++++- src/routing_layer/topo_node_sqlite.h | 4 + tools/chatgui/CMakeLists.txt | 3 +- tools/chatgui/libutun/CMakeLists.txt | 4 +- tools/chatgui/transport/chat_core.c | 16 +- tools/chatgui/transport/chat_core.h | 2 - tools/chatgui/transport/chat_sync.c | 333 +++++++------------------ tools/chatgui/transport/chat_sync.h | 5 - tools/chatgui/transport/utun_node.cpp | 25 +- 16 files changed, 834 insertions(+), 271 deletions(-) create mode 100644 src/routing_layer/topo_group_connect.c create mode 100644 src/routing_layer/topo_group_connect.h diff --git a/src/Makefile.am b/src/Makefile.am index de068c32..ac98f172 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -17,6 +17,7 @@ utun_CORE_SOURCES = \ routing_layer/route_connectivity.c \ routing_layer/conn_mgr.c \ routing_layer/topo_recovery.c \ + routing_layer/topo_group_connect.c \ db_sync.c \ routing_layer/routing.c \ tun_if.c \ @@ -73,6 +74,7 @@ libutun_a_SOURCES = \ routing_layer/route_connectivity.c \ routing_layer/conn_mgr.c \ routing_layer/topo_recovery.c \ + routing_layer/topo_group_connect.c \ db_sync.c \ routing_layer/routing.c \ tun_if.c \ diff --git a/src/routing_layer/conn_mgr.c b/src/routing_layer/conn_mgr.c index 28b24bc7..2e8002ac 100644 --- a/src/routing_layer/conn_mgr.c +++ b/src/routing_layer/conn_mgr.c @@ -54,6 +54,14 @@ static void cm_reverse_timeout_cb(void* arg); static void cm_exchange_timeout_cb(void* arg); static void cm_candidate_ping_timer_cb(void* arg); +static void cm_invite_init_cb(struct ETCP_CONN* conn, void* arg); +static void cm_invite_overall_timeout(void* arg); +static void cm_invite_fail(struct cm_invite_pending* inv, int result); +static void cm_invite_cleanup(struct cm_invite_pending* inv); +static void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn, + const uint8_t* data, size_t len); +static void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len); + struct cm_reverse_pending { struct cm_reverse_pending* next; uint32_t request_id; @@ -72,6 +80,27 @@ struct cm_exchange_pending { uint8_t probes_done; }; +#define CM_INVITE_DEFAULT_TIMEOUT_MS 2000 + +enum CM_INVITE_STATE { + CM_INVITE_CONNECTING = 0, + CM_INVITE_WAIT_INFO = 1, +}; + +struct cm_invite_pending { + struct cm_invite_pending* next; + struct CONN_MGR* mgr; + uint64_t node_id; + uint64_t group_id; + uint32_t timeout_ms; + uint8_t state; + struct ETCP_CONN* conn; + void* overall_timer; + conn_mgr_connect_callback_t cb; + void* cb_arg; + struct TOPO_NODEQ* temp_nq; +}; + struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) { if (!group || !group->instance) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NULL group"); return NULL; } struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR)); @@ -82,7 +111,9 @@ struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) { mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY)); if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "entries alloc failed"); u_free(mgr); return NULL; } mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS; + mgr->invite_list = NULL; etcp_router_bind(mgr->instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler); + etcp_bind(mgr->instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler); 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->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); mgr->bg_ping_cycle_start_tb = get_time_tb(); @@ -96,7 +127,9 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) { if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; } if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; } etcp_router_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR); + etcp_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR); for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]); + while (mgr->invite_list) cm_invite_fail(mgr->invite_list, CONN_MGR_ERR_INTERNAL); struct cm_exchange_pending* ep = mgr->exchange_pending; while (ep) { struct cm_exchange_pending* next = ep->next; @@ -914,6 +947,13 @@ static void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry cm_handle_disconnect(mgr, disc->node_id); } break; + case CONN_MGR_SUBCMD_INVITE_INFO_REQ: + if (mgr && len >= sizeof(struct CONN_MGR_INVITE_INFO_REQ)) + cm_handle_invite_info_req(mgr, conn, dgram, len); + break; + case CONN_MGR_SUBCMD_INVITE_INFO_RESP: + if (mgr) cm_handle_invite_info_resp(mgr, dgram, len); + break; } } @@ -1097,3 +1137,252 @@ static void cm_candidate_ping_timer_cb(void* arg) { mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); } + +/* ═══════════════════════════════════════════════════════════════════════ + * Invite-подключение + * ══════════════════════════════════════════════════════════════════════ */ + +static void sanitize_ch_id(const char* ch_id, char* out, size_t out_sz) { + size_t j = 0; + for (size_t i = 0; ch_id[i] && j < out_sz - 1; i++) + if ((ch_id[i] >= '0' && ch_id[i] <= '9') || (ch_id[i] >= 'a' && ch_id[i] <= 'z') || (ch_id[i] >= 'A' && ch_id[i] <= 'Z') || ch_id[i] == '_') + out[j++] = ch_id[i]; + out[j] = '\0'; +} + +static int cm_check_group_member(struct CONN_MGR* mgr, uint64_t group_id, uint64_t peer_node_id) { + struct TOPO_GROUP* group = topo_groups_find(mgr->instance->topo_groups, group_id); + if (!group) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite_info group 0x%016llx not found", (unsigned long long)group_id); return 0; } + if (group->group_type == TOPO_GROUP_TYPE_CHAT) { + sqlite3* db = mgr->instance->topo_sqlite_db; + if (!db || !group->channel_id[0]) return 0; + char san[64]; sanitize_ch_id(group->channel_id, san, sizeof(san)); + char tbl[80]; snprintf(tbl, sizeof(tbl), "peers_%s", san); + char sql[128]; snprintf(sql, sizeof(sql), "SELECT 1 FROM \"%s\" WHERE node_id=?", tbl); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id); + int is_member = (sqlite3_step(st) == SQLITE_ROW); + sqlite3_finalize(st); + return is_member; + } + return 0; + } + return topo_node_find_by_id(group, peer_node_id) != NULL; +} + +static void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn, + const uint8_t* data, size_t len) { + if (len < sizeof(struct CONN_MGR_INVITE_INFO_REQ)) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ too short (%zu)", len); + return; + } + struct CONN_MGR_INVITE_INFO_REQ* req = (struct CONN_MGR_INVITE_INFO_REQ*)data; + uint64_t peer_node_id = conn->peer_node_id; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ from 0x%016llx group=0x%016llx", + (unsigned long long)peer_node_id, (unsigned long long)req->group_id); + + uint8_t is_member = 0; + if (req->group_id != 0) is_member = (uint8_t)cm_check_group_member(mgr, req->group_id, peer_node_id); + + struct TOPO_GROUP* group = topo_groups_find(mgr->instance->topo_groups, req->group_id); + struct TOPO_GROUP* src_group = group ? group : mgr->group; + const uint8_t* ed_pubkey = src_group->ed25519_public_key; + const char* node_name = (src_group->local_node && src_group->local_node->node) + ? src_group->local_node->node->node_name : NULL; + size_t name_len = node_name ? strlen(node_name) : 0; + if (name_len > 255) name_len = 255; + + size_t resp_size = sizeof(struct CONN_MGR_INVITE_INFO_RESP) + name_len; + struct CONN_MGR_INVITE_INFO_RESP* resp = u_calloc(1, resp_size); + if (!resp) return; + resp->cmd = ETCP_RT_ID_CONN_MGR; + resp->subcmd = CONN_MGR_SUBCMD_INVITE_INFO_RESP; + resp->group_id = req->group_id; + resp->is_member = is_member; + memcpy(resp->ed25519_pubkey, ed_pubkey, SC_PUBKEY_SIZE); + resp->node_name_len = (uint8_t)name_len; + if (name_len) memcpy(resp->node_name, node_name, name_len); + + struct ll_entry* qe = queue_entry_new(0); + if (qe) { qe->dgram = (uint8_t*)resp; qe->len = (uint16_t)resp_size; etcp_send(conn, qe); } + else u_free(resp); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_RESP to 0x%016llx member=%d name_len=%zu", + (unsigned long long)peer_node_id, is_member, name_len); +} + +static void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { + if (len < sizeof(struct CONN_MGR_INVITE_INFO_RESP)) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP too short (%zu)", len); + return; + } + struct CONN_MGR_INVITE_INFO_RESP* resp = (struct CONN_MGR_INVITE_INFO_RESP*)data; + size_t expected = sizeof(struct CONN_MGR_INVITE_INFO_RESP) + (size_t)resp->node_name_len; + if (len < expected) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP truncated (%zu < %zu)", len, expected); + return; + } + + struct cm_invite_pending* prev = NULL; + struct cm_invite_pending* inv = mgr->invite_list; + while (inv) { + if (inv->node_id == resp->group_id ? 0 : 0) { inv = inv->next; continue; } + /* find by state == CM_INVITE_WAIT_INFO and node_id match via conn */ + if (inv->state == CM_INVITE_WAIT_INFO && inv->conn) { + break; + } + prev = inv; inv = inv->next; + } + /* need to find the right invite by peer_node_id from the connection that delivered this */ + struct cm_invite_pending* found = NULL; + for (struct cm_invite_pending* p = mgr->invite_list; p; p = p->next) + if (p->state == CM_INVITE_WAIT_INFO) { found = p; break; } + if (!found) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP — no pending invite in WAIT_INFO state"); return; } + inv = found; + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP node=0x%016llx group=0x%016llx member=%d name_len=%u", + (unsigned long long)inv->node_id, (unsigned long long)resp->group_id, resp->is_member, resp->node_name_len); + + if (inv->group_id != 0 && !resp->is_member) { cm_invite_fail(inv, CONN_MGR_ERR_NOT_MEMBER); return; } + + struct TOPO_NODEQ* nq = inv->temp_nq; + if (nq && nq->node) { + memcpy(nq->node->ed25519_public_key, resp->ed25519_pubkey, SC_PUBKEY_SIZE); + if (nq->node->node_name) { u_free(nq->node->node_name); nq->node->node_name = NULL; } + if (resp->node_name_len) { nq->node->node_name = u_strdup((const char*)resp->node_name); } + sqlite3* db = inv->mgr->instance->topo_sqlite_db; + if (db) topo_node_sqlite_node_put(db, nq, (time_t)(get_time_tb() / 10000)); + } + + if (inv->overall_timer) { uasync_cancel_timeout(inv->mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; } + if (inv->cb) inv->cb(CONN_MGR_OK, inv->node_id, inv->cb_arg); + cm_invite_cleanup(inv); +} + +static void cm_invite_init_cb(struct ETCP_CONN* conn, void* arg) { + struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; + if (!conn || conn->peer_node_id != inv->node_id) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite init cb peer mismatch conn=0x%016llx inv=0x%016llx", + (unsigned long long)(conn ? conn->peer_node_id : 0), (unsigned long long)inv->node_id); + if (conn) cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite INIT OK node=0x%016llx", (unsigned long long)inv->node_id); + inv->state = CM_INVITE_WAIT_INFO; + + struct CONN_MGR_INVITE_INFO_REQ req; + req.cmd = ETCP_RT_ID_CONN_MGR; req.subcmd = CONN_MGR_SUBCMD_INVITE_INFO_REQ; + req.group_id = inv->group_id; + struct ll_entry* qe = queue_entry_new(0); + if (!qe) { cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL); return; } + qe->dgram = u_malloc(sizeof(req)); + if (!qe->dgram) { queue_entry_free(qe); cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL); return; } + memcpy(qe->dgram, &req, sizeof(req)); qe->len = sizeof(req); + etcp_send(conn, qe); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_REQ to 0x%016llx group=0x%016llx", + (unsigned long long)inv->node_id, (unsigned long long)inv->group_id); +} + +static void cm_invite_overall_timeout(void* arg) { + struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; + inv->overall_timer = NULL; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite overall timeout node=0x%016llx", (unsigned long long)inv->node_id); + cm_invite_fail(inv, CONN_MGR_ERR_TIMEOUT); +} + +static void cm_invite_fail(struct cm_invite_pending* inv, int result) { + if (!inv) return; + struct CONN_MGR* mgr = inv->mgr; + if (inv->overall_timer) { uasync_cancel_timeout(mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; } + if (inv->conn) { etcp_connection_close(inv->conn); inv->conn = NULL; } + if (inv->temp_nq) { + struct TOPO_GROUP* group = mgr->group; + queue_remove_data(group->nodes, &inv->temp_nq->ll); + topo_node_free_lists(group, inv->temp_nq); + queue_entry_free(&inv->temp_nq->ll); inv->temp_nq = NULL; + } + struct cm_invite_pending** pp = &mgr->invite_list; + while (*pp) { if (*pp == inv) { *pp = inv->next; break; } pp = &(*pp)->next; } + if (inv->cb) inv->cb(result, inv->node_id, inv->cb_arg); + cm_invite_cleanup(inv); +} + +static void cm_invite_cleanup(struct cm_invite_pending* inv) { + if (!inv) return; + u_free(inv); +} + +int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni, + uint64_t group_id, uint32_t timeout_ms, + 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; } + 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; } + struct TOPO_GROUP* group = mgr->group; + struct TOPO_GROUPS* groups = mgr->instance->topo_groups; + if (!groups) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } + + if (topo_node_find_by_id(group, node_id)) { if (cb) cb(CONN_MGR_ERR_ALREADY_CONNECTED, node_id, cb_arg); return CONN_MGR_ERR_ALREADY_CONNECTED; } + + if (!timeout_ms) timeout_ms = CM_INVITE_DEFAULT_TIMEOUT_MS; + + struct TOPO_NODE* group_ni = u_calloc(1, sizeof(struct TOPO_NODE)); + if (!group_ni) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } + group_ni->node_id = node_id; group_ni->ver = 1; + memcpy(group_ni->public_key, ni->public_key, SC_PUBKEY_SIZE); + for (const struct TOPO_ADDR4* src = ni->v4_addrs; src; src = src->next) { + struct TOPO_ADDR4* a4 = memory_pool_alloc(groups->v4_addr_pool); + if (!a4) continue; + memset(a4, 0, sizeof(*a4)); memcpy(a4->addr, src->addr, 4); a4->port = src->port; + a4->type = TOPO_ADDR_NAT; a4->protocol = TOPO_PROTO_UDP; a4->socket_id = 0; + 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; } + + 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; } + + struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_NODEQ)); + if (!qe) { topo_node_registry_release(groups, group_ni); if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } + struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)qe; + memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry)); + nq->node = group_ni; nq->hash_node_id = node_id; + + struct cm_invite_pending* inv = u_calloc(1, sizeof(struct cm_invite_pending)); + if (!inv) { queue_entry_free(qe); topo_node_registry_release(groups, group_ni); if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; } + inv->mgr = mgr; inv->node_id = node_id; inv->group_id = group_id; + inv->timeout_ms = timeout_ms; inv->state = CM_INVITE_CONNECTING; + inv->temp_nq = nq; inv->cb = cb; inv->cb_arg = cb_arg; + inv->next = mgr->invite_list; mgr->invite_list = inv; + + queue_data_put_with_index(group->nodes, &nq->ll); + + int connected = 0; + for (const struct TOPO_ADDR4* a = group_ni->v4_addrs; a; a = a->next) { + if (a->protocol != TOPO_PROTO_UDP) continue; + struct ETCP_SOCKET* s = mgr->instance->etcp_sockets; + while (s) { + if (s->local_addr.ss_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); + struct sockaddr_storage sa; memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); + struct ETCP_CONN* conn = etcp_connection_create(mgr->instance, NULL); + if (conn) { + etcp_conn_add_init_cbk(conn, cm_invite_init_cb, inv); + sc_init_ctx(&conn->crypto_ctx, &mgr->instance->my_keys); + sc_set_peer_public_key(&conn->crypto_ctx, group_ni->public_key, 0); + if (etcp_link_new(conn, s, &sa, 0)) { inv->conn = conn; connected = 1; break; } + etcp_connection_close(conn); + } + } + s = s->next; + } + if (connected) break; + } + if (!connected) { cm_invite_fail(inv, CONN_MGR_ERR_UNREACHABLE); return CONN_MGR_ERR_UNREACHABLE; } + + inv->overall_timer = uasync_set_timeout(mgr->instance->ua, (int)inv->timeout_ms * 10, inv, cm_invite_overall_timeout, "conn_mgr_invite"); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite started node=0x%016llx group=0x%016llx timeout=%ums", + (unsigned long long)node_id, (unsigned long long)group_id, timeout_ms); + return CONN_MGR_OK; +} diff --git a/src/routing_layer/conn_mgr.h b/src/routing_layer/conn_mgr.h index dc2a4bf8..86e072ef 100644 --- a/src/routing_layer/conn_mgr.h +++ b/src/routing_layer/conn_mgr.h @@ -83,6 +83,7 @@ struct TOPO_GROUP; #define CONN_MGR_ERR_REFUSED -5 ///< Соединение отклонено (зарезервировано) #define CONN_MGR_ERR_ALREADY_CONNECTED -6 ///< Уже подключены #define CONN_MGR_ERR_INTERNAL -7 ///< Внутренняя ошибка +#define CONN_MGR_ERR_NOT_MEMBER -8 ///< Узел не является членом группы /* ==================== Состояния ==================== */ @@ -119,6 +120,8 @@ enum CM_TRY_STATE { #define CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP 0x04 ///< Ответ с нашими кандидатами + замеры RTT до кандидатов инициатора #define CONN_MGR_SUBCMD_INTERM_SELECTED 0x05 ///< Выбранные посредники (финальный шаг INDIRECT-фазы) #define CONN_MGR_SUBCMD_DISCONNECT 0x06 ///< Уведомление о разрыве соединения +#define CONN_MGR_SUBCMD_INVITE_INFO_REQ 0x0A ///< Запрос проверки членства в группе при invite-подключении +#define CONN_MGR_SUBCMD_INVITE_INFO_RESP 0x0B ///< Ответ с членством, ed25519_pubkey, именем узла /** @} */ /** @@ -220,6 +223,38 @@ struct CONN_MGR_DISCONNECT { #pragma pack(pop) +/** + * @struct CONN_MGR_INVITE_INFO_REQ + * @brief Запрос проверки членства в группе (invite-подключение). + * + * Отправляется напрямую по ETCP после INIT-рукопожатия. + * Запрашивает: ed25519_pubkey, имя узла, членство в группе. + */ +#pragma pack(push, 1) +struct CONN_MGR_INVITE_INFO_REQ { + uint8_t cmd; ///< ETCP_RT_ID_CONN_MGR + uint8_t subcmd; ///< CONN_MGR_SUBCMD_INVITE_INFO_REQ + uint64_t group_id; ///< целевая группа для проверки членства +}; +#pragma pack(pop) + +/** + * @struct CONN_MGR_INVITE_INFO_RESP + * @brief Ответ на INVITE_INFO_REQ. + */ +#pragma pack(push, 1) +struct CONN_MGR_INVITE_INFO_RESP { + uint8_t cmd; ///< ETCP_RT_ID_CONN_MGR + uint8_t subcmd; ///< CONN_MGR_SUBCMD_INVITE_INFO_RESP + uint64_t group_id; ///< эхо group_id из запроса + uint8_t is_member; ///< 1 = член группы, 0 = нет + uint8_t reserved[3]; + uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; ///< Ed25519 публичный ключ + uint8_t node_name_len; ///< длина имени (0..255) + uint8_t node_name[0]; ///< UTF-8 имя (variable-length) +}; +#pragma pack(pop) + /* ==================== Размеры пакетов протокола ==================== */ #define CONN_MGR_DIRECT_REQ_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_REQ) ///< Размер заголовка DIRECT_REQ (без addr_count адресов) @@ -228,6 +263,7 @@ struct CONN_MGR_DISCONNECT { #define CONN_MGR_INTERM_EXCHANGE_RESP_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_RESP) ///< Размер INTERM_EXCHANGE_RESP (макс.) #define CONN_MGR_INTERM_SELECTED_SIZE sizeof(struct CONN_MGR_INTERM_SELECTED) ///< Размер INTERM_SELECTED #define CONN_MGR_DISCONNECT_SIZE sizeof(struct CONN_MGR_DISCONNECT) ///< Размер DISCONNECT +#define CONN_MGR_INVITE_INFO_REQ_SIZE sizeof(struct CONN_MGR_INVITE_INFO_REQ) ///< Размер INVITE_INFO_REQ /* ==================== Callback-типы и внутренние структуры ==================== */ @@ -319,6 +355,7 @@ struct CONN_MGR { struct cm_reverse_pending* reverse_pending; ///< Связный список ожидающих REVERSE-подключений struct cm_exchange_pending* exchange_pending; ///< Связный список ожидающих INDIRECT-обменов + struct cm_invite_pending* invite_list; ///< Связный список ожидающих invite-подключений }; /* ==================== Жизненный цикл ==================== */ @@ -411,6 +448,24 @@ int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t t */ void conn_mgr_set_direct_timeout_ms(struct CONN_MGR* mgr, uint32_t timeout_ms); +/** + * @brief Подключается к новому узлу по invite-данным и проверяет членство в группе. + * @param mgr менеджер группы + * @param ni данные узла (caller-owned, только для чтения): node_id, public_key, v4_addrs + * @param group_id группа для проверки членства (0 = без проверки) + * @param timeout_ms общий таймаут (0 = 2000ms по умолчанию) + * @param cb callback результата + * @param cb_arg аргумент для callback + * @return CONN_MGR_OK если процесс запущен, CONN_MGR_ERR_* при ошибке + * + * Поток: DIRECT-подключение → INIT → INVITE_INFO_REQ → INVITE_INFO_RESP → успех/ошибка. + * При успехе узел добавляется в group->nodes и сохраняется в SQLite. + * Callback гарантированно вызывается (успех, таймаут, ошибка членства). + */ +int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni, + uint64_t group_id, uint32_t timeout_ms, + conn_mgr_connect_callback_t cb, void* cb_arg); + /** * @brief Получает текущий статус подключения к ноде. * @param mgr менеджер diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index b7ada200..14039b8b 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -23,6 +23,7 @@ #include "route_connectivity.h" #include "control_server.h" #include "conn_mgr.h" +#include "topo_group_connect.h" #include "topo_recovery.h" @@ -238,6 +239,7 @@ static void topo_group_destroy(struct TOPO_GROUP* group) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "group_id=%016llx", (unsigned long long)group->group_id); topo_recovery_cancel_all(group); + topo_group_connect_destroy(group); struct ll_entry* e; while ((e = queue_data_get(group->senders_list)) != NULL) queue_entry_free(e); @@ -381,6 +383,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou } DEBUG_INFO(DEBUG_CATEGORY_BGP, "Group created: group_id=%016llx type=%u ch_id=%s", (unsigned long long)group_id, group_type, group->channel_id); + if (group_type == TOPO_GROUP_TYPE_CHAT) topo_group_connect_init(group); return group; } @@ -406,6 +409,7 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id); topo_group_send_table_request(group, conn); + topo_group_connect_on_up(group, conn); } void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { @@ -451,6 +455,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { if (cascaded > 0) topo_recovery_start(group); + topo_group_connect_on_down(group, conn); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d", conn->log_name, nodes_removed); } diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 9d83300f..5af521d8 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -2,6 +2,12 @@ * @file topo_group.h * @brief BGP-подобный обмен топологией узлов между пирами через ETCP. * + * Идея: модуль автоматически выстраивает карту маршрутизации между узлами, + * используя только пассивное наблюдение (подписка на события) + * Можно иметь много групп узлов. Группа - это список узлов. + * Для каждой группы выстраивается своя независимая таблица маршрутизации. + * Один узел может входить в любое число групп. + * Механика: * - Узлы обмениваются информацией друг о друге (pubkey, адреса, подсети) * через NODEINFO/WITHDRAW сообщения, маршрутизируемые по ETCP. @@ -43,6 +49,7 @@ struct UTUN_INSTANCE; struct NAT_DETECTION; struct CONN_MGR; struct TOPO_RECOVERY_CTX; +struct TOPO_GROUP_CONNECT; /** Callback when a node's info (pubkeys, addresses) is persisted in the DB */ typedef void (*topo_node_updated_fn)(struct UTUN_INSTANCE* inst, uint64_t node_id, @@ -132,6 +139,7 @@ struct TOPO_GROUP { char channel_id[64]; // channel_id для групп типа CHAT struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления + struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT) }; /** diff --git a/src/routing_layer/topo_group_connect.c b/src/routing_layer/topo_group_connect.c new file mode 100644 index 00000000..923e702e --- /dev/null +++ b/src/routing_layer/topo_group_connect.c @@ -0,0 +1,260 @@ +/* + * topo_group_connect.c — авто-подключение к узлам группы при старте + * + * Phase 1: одновременный запуск conn_mgr_connect_node для пиров с connected=1. + * Таймаут = direct_timeout_ms. По таймауту отменяем незавершённые и помечаем + * connected=0 в БД. Если хоть один подключился → done. + * + * Phase 2: последовательный перебор пиров с публичными адресами до первого успеха. + * + * Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями. + * UP: инкремент только для первого соединения к peer_node_id. + * DOWN: декремент только если нет других conn к peer_node_id И нет indirect-путей в nodeinfo. + * + * restart: перезапуск авто-подключения при phase=done и active_conn_count==0. + */ + +#include "topo_group_connect.h" +#include "topo_group.h" +#include "topo_node.h" +#include "topo_node_sqlite.h" +#include "conn_mgr.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "etcp.h" + +/* ─── внутренние константы ─── */ +#define TGC_ID "topo_group_connect" +#define TGC_PHASE_ONE 0 +#define TGC_PHASE_TWO 1 +#define TGC_PHASE_DONE 2 + +struct TOPO_GROUP_CONNECT { + struct TOPO_GROUP* group; + void* phase_timer; + uint8_t phase; + uint8_t active; + int pending; + int connected_count; + int active_conn_count; // число уникальных peer-узлов с живыми соединениями + int cursor; + uint64_t* candidate_ids; + int candidate_count; +}; + +static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc); +static void tgc_phase2_result(int result, uint64_t node_id, void* arg); +static void tgc_phase1_timeout(void* arg); +static void tgc_phase1_result(int result, uint64_t node_id, void* arg); + +/* ═══════════════════════════════════════════════════════════════════════ + * Жизненный цикл + * ══════════════════════════════════════════════════════════════════════ */ + +int topo_group_connect_init(struct TOPO_GROUP* group) { + if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || group->connect) + return -1; + struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc)); + if (!gc) return -1; + gc->group = group; gc->active = 1; + group->connect = gc; + + uint64_t* ids = NULL; int count = 0; + sqlite3* db = group->instance->topo_sqlite_db; + if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) { + if (ids) { u_free(ids); ids = NULL; } + gc->phase = TGC_PHASE_TWO; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s no connected peers → Phase 2", TGC_ID, group->channel_id); + tgc_phase2_try_next(gc); + return 0; + } + + gc->phase = TGC_PHASE_ONE; gc->pending = count; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s Phase 1 launching %d connects", TGC_ID, group->channel_id, count); + for (int i = 0; i < count; i++) + conn_mgr_connect_node(group->conn_mgr, ids[i], 0, tgc_phase1_result, gc); + u_free(ids); + + gc->phase_timer = uasync_set_timeout(group->instance->ua, + (int)group->conn_mgr->direct_timeout_ms * 10, gc, tgc_phase1_timeout, "tgc_phase1"); + return 0; +} + +void topo_group_connect_destroy(struct TOPO_GROUP* group) { + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc) return; + group->connect = NULL; + if (gc->phase_timer) { uasync_cancel_timeout(group->instance->ua, gc->phase_timer); gc->phase_timer = NULL; } + u_free(gc->candidate_ids); + u_free(gc); +} + +void topo_group_connect_restart(struct TOPO_GROUP* group) { + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc) { topo_group_connect_init(group); return; } + if (gc->phase != TGC_PHASE_DONE) return; + if (gc->active_conn_count > 0) return; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s", TGC_ID, group->channel_id); + topo_group_connect_destroy(group); + topo_group_connect_init(group); +} + +int topo_group_connect_active_count(struct TOPO_GROUP* group) { + struct TOPO_GROUP_CONNECT* gc = group->connect; + return gc ? gc->active_conn_count : 0; +} + +/* ═══════════════════════════════════════════════════════════════════════ + * ON UP / DOWN + * ══════════════════════════════════════════════════════════════════════ */ + +void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return; + uint64_t peer = conn->peer_node_id; + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc) return; + + /* уже есть другое прямое соединение к этому peer? */ + struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; + while (e) { + struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) { + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx already connected via other conn, skip count", + TGC_ID, (unsigned long long)peer); + return; + } + e = e->next; + } + + gc->active_conn_count++; + sqlite3* db = group->instance->topo_sqlite_db; + if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s active=%d", + TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); +} + +void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return; + uint64_t peer = conn->peer_node_id; + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc) return; + + /* есть ли другие прямые соединения к этому peer? */ + struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; + while (e) { + struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) { + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx still has other direct conn, skip", + TGC_ID, (unsigned long long)peer); + return; + } + e = e->next; + } + + /* есть ли indirect-пути в nodeinfo? */ + struct TOPO_NODEQ* nq = topo_node_find_by_id(group, peer); + if (nq && nq->paths) { + struct ll_entry* pe = nq->paths->head; + while (pe) { + struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe; + if (path->conn && path->conn != conn) { + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx has indirect path, skip", + TGC_ID, (unsigned long long)peer); + return; + } + pe = pe->next; + } + } + + gc->active_conn_count--; + sqlite3* db = group->instance->topo_sqlite_db; + if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d", + TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); +} + +/* ═══════════════════════════════════════════════════════════════════════ + * Phase 1 + * ══════════════════════════════════════════════════════════════════════ */ + +static void tgc_phase1_result(int result, uint64_t node_id, void* arg) { + struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; + if (gc->phase != TGC_PHASE_ONE) return; + gc->pending--; + if (result == CONN_MGR_OK) gc->connected_count++; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 result node=0x%016llx %s pending=%d connected=%d", + TGC_ID, (unsigned long long)node_id, result == CONN_MGR_OK ? "OK" : "FAIL", gc->pending, gc->connected_count); +} + +static void tgc_phase1_timeout(void* arg) { + struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; + gc->phase_timer = NULL; + if (gc->phase != TGC_PHASE_ONE) return; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 timeout connected=%d pending=%d ch=%s", + TGC_ID, gc->connected_count, gc->pending, gc->group->channel_id); + + sqlite3* db = gc->group->instance->topo_sqlite_db; + uint64_t* ids = NULL; int count = 0; + topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count); + for (int i = 0; i < count; i++) { + uint8_t status; uint8_t conn_type; + conn_mgr_get_status(gc->group->conn_mgr, ids[i], &status, &conn_type); + if (status != CONN_MGR_STATE_CONNECTED) { + conn_mgr_cancel_callback(gc->group->conn_mgr, ids[i], tgc_phase1_result, gc); + topo_node_sqlite_set_connected(db, gc->group->channel_id, ids[i], 0); + } + } + u_free(ids); + + if (gc->connected_count > 0) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — %d connected, skipping Phase 2", TGC_ID, gc->connected_count); + gc->phase = TGC_PHASE_DONE; + return; + } + gc->phase = TGC_PHASE_TWO; + tgc_phase2_try_next(gc); +} + +/* ═══════════════════════════════════════════════════════════════════════ + * Phase 2 + * ══════════════════════════════════════════════════════════════════════ */ + +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->cursor = 0; + DEBUG_INFO(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++]; + uint8_t status; uint8_t conn_type; + conn_mgr_get_status(gc->group->conn_mgr, nid, &status, &conn_type); + if (status == CONN_MGR_STATE_CONNECTED) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 skip 0x%016llx already connected", TGC_ID, (unsigned long long)nid); + continue; + } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d)", TGC_ID, + (unsigned long long)nid, gc->cursor, gc->candidate_count); + conn_mgr_connect_node(gc->group->conn_mgr, nid, 0, tgc_phase2_result, gc); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id); + gc->phase = TGC_PHASE_DONE; +} + +static void tgc_phase2_result(int result, uint64_t node_id, void* arg) { + struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; + if (gc->phase != TGC_PHASE_TWO) return; + if (result == CONN_MGR_OK) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — done", TGC_ID, (unsigned long long)node_id); + gc->phase = TGC_PHASE_DONE; + return; + } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — try next", TGC_ID, (unsigned long long)node_id); + tgc_phase2_try_next(gc); +} diff --git a/src/routing_layer/topo_group_connect.h b/src/routing_layer/topo_group_connect.h new file mode 100644 index 00000000..86f88503 --- /dev/null +++ b/src/routing_layer/topo_group_connect.h @@ -0,0 +1,30 @@ +/** + * @file topo_group_connect.h + * @brief Авто-подключение к узлам группы при старте и отслеживание состояния connected в БД. + * + * При старте CHAT-группы: + * Phase 1 — одновременный запуск подключений ко всем пирам с connected=1. + * Если хоть один подключился → done. + * Phase 2 — последовательный перебор по пирам с публичными адресами. + * До первого успешного подключения. + * + * Отслеживание: на UP → peers_.connected=1, на DOWN (последний conn) → 0. + * active_conn_count — количество уникальных peer-узлов с живыми соединениями. + * restart — перезапуск авто-подключения (если phase=done и нет подключений). + */ +#ifndef TOPO_GROUP_CONNECT_H +#define TOPO_GROUP_CONNECT_H + +#include + +struct TOPO_GROUP; +struct ETCP_CONN; + +int topo_group_connect_init(struct TOPO_GROUP* group); +void topo_group_connect_destroy(struct TOPO_GROUP* group); +void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn); +void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn); +int topo_group_connect_active_count(struct TOPO_GROUP* group); +void topo_group_connect_restart(struct TOPO_GROUP* group); + +#endif diff --git a/src/routing_layer/topo_node_sqlite.c b/src/routing_layer/topo_node_sqlite.c index 7d3729c9..db7b88a3 100644 --- a/src/routing_layer/topo_node_sqlite.c +++ b/src/routing_layer/topo_node_sqlite.c @@ -92,6 +92,9 @@ int topo_node_sqlite_init(sqlite3* db) { snprintf(isql, sizeof(isql), "CREATE INDEX IF NOT EXISTS \"idx_%s_type_rtt\" ON \"%s\" (node_type, node_RTT)", tbl, tbl); sqlite3_exec(db, isql, NULL, NULL, NULL); + snprintf(isql, sizeof(isql), + "ALTER TABLE \"%s\" ADD COLUMN connected INTEGER NOT NULL DEFAULT 0", tbl); + sqlite3_exec(db, isql, NULL, NULL, NULL); } sqlite3_finalize(s); } @@ -274,7 +277,8 @@ int topo_node_sqlite_channel_put(sqlite3* db, const char* channel_id, " node_type INTEGER NOT NULL DEFAULT 0," /* вычисляемый тип узла (topo_node_sqlite_nodeinfo_updated): 0=неизвестно 1=прямой(addr_type NETIF/DIRECT/NAT_STRICT) 2=EIM NAT 3=(не исп.) 4=суперузел(DIRECT+"supernode=yes") */ - " node_RTT INTEGER" /* RTT до узла в мс, часть индекса (node_type, node_RTT) */ + " node_RTT INTEGER," /* RTT до узла в мс, часть индекса (node_type, node_RTT) */ + " connected INTEGER NOT NULL DEFAULT 0" /* 1 = было активное подключение к узлу */ ")"; int ddl_sz = snprintf(NULL, 0, ddl_tmpl, peers_tbl); @@ -724,3 +728,60 @@ struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* gr (unsigned long long)node_id, name ? name : "", topo_list_count((struct _topo_head*)v4_head)); return ni; } + +int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected) { + if (!db || !channel_id) return -1; + char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl)); + char sql[128]; snprintf(sql, sizeof(sql), "UPDATE \"%s\" SET connected=? WHERE node_id=?", peers_tbl); + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1; + sqlite3_bind_int(st, 1, connected ? 1 : 0); + sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id); + int rc = sqlite3_step(st); sqlite3_finalize(st); + return (rc == SQLITE_DONE) ? 0 : -1; +} + +int topo_node_sqlite_get_connected_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[128]; snprintf(sql, sizeof(sql), "SELECT node_id FROM \"%s\" WHERE connected=1", 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_public_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[280]; + snprintf(sql, sizeof(sql), + "SELECT p.node_id FROM \"%s\" p" + " WHERE EXISTS (SELECT 1 FROM node_addresses a" + " WHERE a.node_id=p.node_id AND a.family=4 AND a.addr_type!=0)" + " ORDER BY p.node_RTT", 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; +} diff --git a/src/routing_layer/topo_node_sqlite.h b/src/routing_layer/topo_node_sqlite.h index cda6d134..fe606eec 100644 --- a/src/routing_layer/topo_node_sqlite.h +++ b/src/routing_layer/topo_node_sqlite.h @@ -67,4 +67,8 @@ void topo_node_sqlite_update_rtt(sqlite3* db, uint64_t node_id, uint16_t rtt); */ struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* groups, uint64_t node_id); +int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected); +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); + #endif diff --git a/tools/chatgui/CMakeLists.txt b/tools/chatgui/CMakeLists.txt index 487691fd..ba632574 100644 --- a/tools/chatgui/CMakeLists.txt +++ b/tools/chatgui/CMakeLists.txt @@ -80,7 +80,6 @@ add_executable(chatgui transport/config_updater.cpp transport/gui_bridge_impl.cpp transport/chat_core.c - transport/chat_conn_mgr.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c @@ -92,7 +91,7 @@ add_executable(chatgui target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE) -set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_conn_mgr.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/miniaudio_impl.c PROPERTIES LANGUAGE C) +set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/miniaudio_impl.c PROPERTIES LANGUAGE C) if(WIN32) target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread) else() diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index 7f1f9cc2..9e935040 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -67,9 +67,7 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/db_sync.c ${SRC_DIR}/routing_layer/conn_mgr.c ${SRC_DIR}/routing_layer/topo_recovery.c - ${TRANSPORT_DIR}/chat_sync.c - ${TRANSPORT_DIR}/merkle_sync.c - ${TRANSPORT_DIR}/member_sync.c + ${SRC_DIR}/routing_layer/topo_group_connect.c ${SRC_DIR}/routing_layer/routing.c ${SRC_DIR}/tun_if.c ${SRC_DIR}/tun_route.c diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 00771753..2ef16763 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -6,11 +6,11 @@ */ #include "chat_core.h" -#include "chat_conn_mgr.h" #include "db_sync.h" #include "gui_bridge.h" #include "../../../lib/json_flat.h" #include "topo_node_sqlite.h" +#include "topo_group.h" #include "member_sync.h" #include "../../../src/utun_instance.h" @@ -39,9 +39,8 @@ static struct chat_core_ctx { struct UTUN_INSTANCE* inst; sqlite3* db; - uint8_t shared_db; /* db == inst->topo_sqlite_db, do not close */ + uint8_t shared_db; uint64_t my_node_id; - struct chat_conn_mgr* conn_mgr; /* db_sync instances per channel */ struct DB_SYNC_INSTANCE** si; @@ -136,7 +135,6 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) { } } - g_cc.conn_mgr = chat_conn_mgr_init(inst, g_cc.db); g_cc.initialized = 1; chat_core_sync_my_addresses(); @@ -237,16 +235,11 @@ int chat_core_is_initialized(void) { return g_cc.initialized; } -struct chat_conn_mgr* chat_conn_mgr_get(void) { - return g_cc.conn_mgr; -} - void chat_core_destroy(struct UTUN_INSTANCE* inst) { (void)inst; if (!g_cc.initialized) return; g_cc.initialized = 0; - chat_conn_mgr_destroy(g_cc.conn_mgr); g_cc.conn_mgr = NULL; if (g_cc.db && !g_cc.shared_db) { sqlite3_close(g_cc.db); } g_cc.db = NULL; g_cc.inst = NULL; DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID); @@ -670,6 +663,11 @@ void chat_core_ensure_channel_ready(const char* ch_id) { uint64_t ch_hash = 0; { const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); } struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash); + + /* create TOPO_GROUP_TYPE_CHAT for auto-connect */ + uint64_t gid = strtoull(ch_id, NULL, 10); + if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid)) + topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id); if (si) { si_register(si, ch_id); db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id)); diff --git a/tools/chatgui/transport/chat_core.h b/tools/chatgui/transport/chat_core.h index dc45860a..bcbf93e0 100644 --- a/tools/chatgui/transport/chat_core.h +++ b/tools/chatgui/transport/chat_core.h @@ -23,8 +23,6 @@ struct sqlite3* chat_core_get_db(void); struct UTUN_INSTANCE* chat_core_get_inst(void); int chat_core_is_initialized(void); -struct chat_conn_mgr* chat_conn_mgr_get(void); - /* ── Отправка сообщения (GUI → uasync) ── */ struct chat_msg_submit { diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index fa32f0dc..2f8c2c17 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -1,8 +1,10 @@ #include "chat_sync.h" #include "chat_core.h" -#include "chat_conn_mgr.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" +#include "topo_group.h" +#include "conn_mgr.h" +#include "topo_group_connect.h" #include "../../../lib/json_flat.h" #include "member_sync.h" #include "merkle_sync.h" @@ -24,238 +26,6 @@ static struct chat_sync* g_cs = NULL; -/* ═══════════════════════════════════════════════════════════════════════ - * Auto-connect: cursor-based peer iteration with GC every 1s - * ══════════════════════════════════════════════════════════════════════ */ - -#define AC_MAX_FLIGHTS 10 -#define AC_RETRY_MS 1000 -#define AC_GC_TIMEOUT_MS 3000 -#define AC_ID "auto_connect" - -struct ac_flight { - uint64_t node_id; - uint8_t active; /* 1 = connection attempt in progress */ - uint64_t created_tb; -}; - -struct auto_connect { - struct UTUN_INSTANCE* inst; - void* retry_timer; - uint8_t active; - - char** channel_ids; /* loaded once, refreshed on cursor wrap */ - int channel_count; - int ch_cursor; /* current channel index */ - int peer_cursor; /* current peer index within channel */ - - struct ac_flight flights[AC_MAX_FLIGHTS]; -}; - -static struct auto_connect* g_ac = NULL; - -static void ac_retry_timer_cb(void* arg); -static void ac_result_cb(int result, uint64_t node_id, void* arg); - -/* ── load channel IDs into auto_connect ── */ - -static int ac_load_channels(struct auto_connect* ac) { - if (ac->channel_ids) { - for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]); - u_free(ac->channel_ids); - ac->channel_ids = NULL; - } - ac->channel_count = 0; - uint8_t buf[4096]; size_t buf_len; - if (chat_core_list_channels(buf, sizeof(buf), &buf_len) != 0 || buf_len < 2) return 0; - uint16_t cnt; memcpy(&cnt, buf, 2); - const uint8_t* p = buf + 2; size_t rem = buf_len - 2; - ac->channel_ids = u_calloc(cnt, sizeof(char*)); - if (!ac->channel_ids) return 0; - for (uint16_t i = 0; i < cnt && rem >= 1; i++) { - uint8_t id_len = *p++; rem--; - if (rem < id_len) break; - ac->channel_ids[i] = u_malloc(id_len + 1); - if (ac->channel_ids[i]) { memcpy(ac->channel_ids[i], p, id_len); ac->channel_ids[i][id_len] = '\0'; ac->channel_count++; } - p += id_len; rem -= id_len; - } - return ac->channel_count; -} - -/* ── GC: close expired flights (no link_status after AC_GC_TIMEOUT_MS) ── */ - -static void ac_gc(struct auto_connect* ac) { - uint64_t now = get_time_tb(); - uint64_t deadline = (uint64_t)AC_GC_TIMEOUT_MS * 10; - for (int i = 0; i < AC_MAX_FLIGHTS; i++) { - if (!ac->flights[i].active) continue; - if (now - ac->flights[i].created_tb < deadline) continue; - uint64_t nid = ac->flights[i].node_id; - /* check if link already UP — if so, just free slot (already connected) */ - int found_up = 0; - { - struct ll_entry* entry = ac->inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (ce->conn->peer_node_id == nid) { - struct ETCP_LINK* l = ce->conn->links; - while (l) { if (l->initialized && l->link_status) { found_up = 1; break; } l = l->next; } - if (found_up) break; - } - entry = entry->next; - } - } - if (found_up) { - chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id); - memset(&ac->flights[i], 0, sizeof(ac->flights[i])); - continue; - } - /* expired and not up — cancel */ - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: GC closing flight node=0x%016llx (expired)", AC_ID, (unsigned long long)nid); - chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id); - memset(&ac->flights[i], 0, sizeof(ac->flights[i])); - } -} - -/* ── count occupied flight slots ── */ - -static int ac_flight_count(struct auto_connect* ac) { - int n = 0; - for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (ac->flights[i].active) n++; - return n; -} - -/* ── find first free flight slot ── */ - -static int ac_find_free_slot(struct auto_connect* ac) { - for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (!ac->flights[i].active) return i; - return -1; -} - -/* ── check if node already has active ETCP connection ── */ - -static int ac_node_has_conn(struct UTUN_INSTANCE* inst, uint64_t nid) { - if (!inst->connections) return 0; - struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&nid); - return e != NULL; -} - -/* ── fill up to AC_MAX_FLIGHTS by advancing cursor ── */ - -static void ac_fill(struct auto_connect* ac) { - if (ac->channel_count == 0) return; - /* try all channels up to 2 full cycles without finding a candidate → give up */ - int max_iter = ac->channel_count * 2 + 10; - while (ac_flight_count(ac) < AC_MAX_FLIGHTS && max_iter-- > 0) { - int slot = ac_find_free_slot(ac); - if (slot < 0) break; - char* ch_id = ac->channel_ids[ac->ch_cursor]; - uint8_t pb[2048]; size_t plen; - if (chat_core_list_peers(ch_id, pb, sizeof(pb), &plen) != 0 || plen < 2) { - /* channel has no peers or error — skip to next */ - ac->ch_cursor++; ac->peer_cursor = 0; - if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); } - continue; - } - uint16_t pc; memcpy(&pc, pb, 2); - if (ac->peer_cursor >= (int)pc) { - ac->ch_cursor++; ac->peer_cursor = 0; - if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); } - continue; - } - /* get next peer */ - uint64_t nid; memcpy(&nid, pb + 2 + ac->peer_cursor * 8, 8); - ac->peer_cursor++; - if (nid == 0 || nid == ac->inst->node_id) continue; - - /* already connected? */ - if (ac_node_has_conn(ac->inst, nid)) continue; - - /* already in a flight slot? */ - int dup = 0; - for (int i = 0; i < AC_MAX_FLIGHTS; i++) - if (ac->flights[i].active && ac->flights[i].node_id == nid) { dup = 1; break; } - if (dup) continue; - - /* launch */ - struct ac_flight* f = &ac->flights[slot]; - f->node_id = nid; - f->created_tb = get_time_tb(); - chat_conn_mgr_connect(chat_conn_mgr_get(), nid, ac_result_cb, f); f->active = 1; - if (f->node_id != nid) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: flight to 0x%016llx FAIL (ca_state=NULL, sync fail in chat_conn_mgr_connect_auto)", AC_ID, (unsigned long long)nid); - f->node_id = 0; f->created_tb = 0; - continue; - } - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: flight[%d] launched node=0x%016llx", AC_ID, slot, (unsigned long long)nid); - uint8_t evt[7]; evt[0] = 0; uint16_t t = 0, n = (uint16_t)ac->channel_count, s = 0; - memcpy(evt + 1, &t, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &s, 2); - gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); - } -} - -/* ── result callback (fired by chat_conn_mgr_connect_auto) ── */ - -static void ac_result_cb(int result, uint64_t node_id, void* arg) { - struct ac_flight* f = (struct ac_flight*)arg; - if (!g_ac || !g_ac->active) return; - const char* rs = result == CC_OK ? "OK" : "FAIL"; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: result node=0x%016llx %s", AC_ID, (unsigned long long)node_id, rs); - f->active = 0; f->node_id = 0; f->created_tb = 0; -} - -/* ── retry timer callback (every AC_RETRY_MS) ── */ - -static void ac_retry_timer_cb(void* arg) { - struct auto_connect* ac = (struct auto_connect*)arg; - ac->retry_timer = NULL; - if (!ac->active) return; - if (ac->channel_count == 0) ac_load_channels(ac); - ac_gc(ac); - ac_fill(ac); - ac->retry_timer = uasync_set_timeout(ac->inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry"); -} - -/* ═══════════════════════════════════════════════════════════════════════ - * Public auto_connect API - * ══════════════════════════════════════════════════════════════════════ */ - -void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst) { - if (!inst || !inst->ua) return; - chat_sync_auto_connect_stop(); - struct auto_connect* ac = u_calloc(1, sizeof(*ac)); - if (!ac) return; - ac->inst = inst; - ac->active = 1; - g_ac = ac; - ac_load_channels(ac); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: start, %d channels loaded", AC_ID, ac->channel_count); - ac_gc(ac); - ac_fill(ac); - ac->retry_timer = uasync_set_timeout(inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry"); -} - -void chat_sync_auto_connect_stop(void) { - struct auto_connect* ac = g_ac; - if (!ac) return; - ac->active = 0; g_ac = NULL; - if (ac->retry_timer) { uasync_cancel_timeout(ac->inst->ua, ac->retry_timer); ac->retry_timer = NULL; } - for (int i = 0; i < AC_MAX_FLIGHTS; i++) - if (ac->flights[i].active) { chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id); memset(&ac->flights[i], 0, sizeof(ac->flights[i])); } - if (ac->channel_ids) { for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]); u_free(ac->channel_ids); } - uint8_t evt[7]; evt[0] = 2; uint16_t z = 0; - memcpy(evt + 1, &z, 2); memcpy(evt + 3, &z, 2); memcpy(evt + 5, &z, 2); - gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); - u_free(ac); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: stopped", AC_ID); -} - -void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: switch group 0x%016llx", AC_ID, (unsigned long long)new_group_id); - chat_sync_auto_connect_stop(); - chat_sync_auto_connect_start(inst); -} - struct channel_cache { char channel_id[64]; uint64_t* peer_ids; @@ -399,7 +169,7 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { case CS_MSG_ERROR: { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id); if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; } - uint8_t err[20]; memcpy(err, &peer, 8); int r = CC_ERR_NOT_FOUND; memcpy(err + 8, &r, 4); + uint8_t err[20]; memcpy(err, &peer, 8); int r = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4); memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0; break; @@ -420,7 +190,7 @@ static void cs_info_req_timeout_cb(void* arg) { CS_ID, (unsigned long long)cs->pending_invite_node_id, (unsigned long long)cs->pending_invite_ch_id); uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8); - int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4); + int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4); memcpy(err + 12, &cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); cs->pending_invite_node_id = 0; @@ -435,7 +205,7 @@ static void cs_join_timeout_cb(void* arg) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu", CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id); uint8_t err[20]; memcpy(err, &peer, 8); - int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4); + int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4); memcpy(err + 12, &cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); cs->pending_invite_node_id = 0; @@ -785,7 +555,6 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { struct chat_sync* cs = g_cs; if (!cs || !inst) return; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: destroy", CS_ID); - chat_sync_auto_connect_stop(); if (cs->sync_timer) { uasync_cancel_timeout(inst->ua, cs->sync_timer); cs->sync_timer = NULL; } if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); } member_sync_destroy(inst); @@ -803,15 +572,87 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { } void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { - (void)inst; - chat_conn_mgr_connect(chat_conn_mgr_get(), node_id, NULL, NULL); + (void)inst; (void)node_id; } -struct cm_invite_wrap { struct chat_conn_mgr* mgr; struct chat_invite inv; }; +struct chat_invite { + uint64_t channel_id; + uint64_t node_id; + uint8_t pubkey[32]; + uint8_t* addrs_data; + int addr_count; +}; + +struct cm_invite_wrap { struct chat_invite inv; }; static void cm_invite_trampoline(void* arg) { struct cm_invite_wrap* w = (struct cm_invite_wrap*)arg; - chat_conn_mgr_connect_from_invite(w->mgr, &w->inv); + struct chat_invite* inv = &w->inv; + struct UTUN_INSTANCE* inst = chat_core_get_inst(); + sqlite3* db = chat_core_get_db(); + uint64_t node_id = inv->node_id; + uint64_t channel_id = inv->channel_id; + + /* save pubkey to nodes */ + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(db, "INSERT INTO nodes(node_id,x25519_pubkey,created_at) VALUES(?,?,?)" + " ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey," + " created_at=COALESCE(nodes.created_at, excluded.created_at)", -1, &st, NULL); + if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 3, 0); sqlite3_step(st); sqlite3_finalize(st); } + + /* save invite addresses (socket_id=0) */ + sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", -1, &st, NULL); + if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_step(st); sqlite3_finalize(st); } + sqlite3_prepare_v2(db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id) VALUES(?,?,1,?,?,0,?)", -1, &st, NULL); + if (st) { + const uint8_t* ap = inv->addrs_data; + for (int i = 0; i < inv->addr_count; i++) { + uint8_t family = *ap++; ap++; /* skip sock_id */ + if (family != 4) { ap += 18; continue; } + sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_int(st, 2, 4); + sqlite3_bind_blob(st, 3, ap, 4, SQLITE_STATIC); ap += 4; + uint16_t port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; + sqlite3_bind_int(st, 4, (int)port); sqlite3_bind_int(st, 5, 1); /* addr_type=1 (DIRECT) */ + sqlite3_step(st); sqlite3_reset(st); + } + sqlite3_finalize(st); + } + + char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)channel_id); + chat_core_ensure_channel_ready(ch_str); + + /* ensure TOPO_GROUP_TYPE_CHAT group exists */ + if (inst->topo_groups) { + uint64_t gid = channel_id; + if (!topo_groups_find(inst->topo_groups, gid)) + topo_groups_create_group(inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_str); + } + + /* build TOPO_NODE from invite data */ + struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); + if (!ni) { u_free(w); return; } + ni->node_id = node_id; memcpy(ni->public_key, inv->pubkey, 32); + const uint8_t* ap = inv->addrs_data; + for (int i = 0; i < inv->addr_count; i++) { + uint8_t family = *ap++; ap++; /* skip sock_id */ + if (family != 4) { ap += 18; continue; } + struct TOPO_ADDR4* a4 = u_calloc(1, sizeof(*a4)); + if (!a4) break; + memcpy(a4->addr, ap, 4); ap += 4; + a4->port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; + a4->type = TOPO_ADDR_NAT; a4->protocol = TOPO_PROTO_UDP; + a4->next = ni->v4_addrs; ni->v4_addrs = a4; + } + + if (inst->topo_groups && ni->v4_addrs) { + struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id); + if (g && g->conn_mgr) + conn_mgr_connect_from_invite(g->conn_mgr, ni, channel_id, 0, NULL, NULL); + } + + /* free caller-owned ni */ + while (ni->v4_addrs) { struct TOPO_ADDR4* n = ni->v4_addrs->next; u_free(ni->v4_addrs); ni->v4_addrs = n; } + u_free(ni); u_free(w); } @@ -846,8 +687,8 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start ch=%llu node=0x%016llx pubkey=%016llx... addrs=%d", CS_ID, channel_id, node_id, *(const uint64_t*)pubkey_bin, addr_count); - struct cm_invite_wrap { struct chat_conn_mgr* mgr; struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap)); - w->mgr = chat_conn_mgr_get(); w->inv = *inv; u_free(inv); + struct cm_invite_wrap { struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap)); + w->inv = *inv; u_free(inv); gui_bridge_post_uasync_fn( (void(*)(void*))cm_invite_trampoline, w); } @@ -1193,7 +1034,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer, if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(joiner) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc); { char j_name[128]; json_flat_get(joiner_userinfo, "name", j_name, sizeof(j_name)); topo_node_sqlite_node_update_verified(db, node_id, j_name, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } - chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id); + /* save/update node addresses */ DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN save addrs node=0x%016llx addr_cnt=%d", CS_ID, (unsigned long long)node_id, addr_cnt); @@ -1348,7 +1189,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, { char pn[256]; json_flat_get(peer_userinfo, "name", pn, sizeof(pn)); topo_node_sqlite_node_update_verified(db, node_id, pn, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } - chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] WELCOME save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); for (uint8_t j = 0; j < ac; j++) { @@ -1451,7 +1292,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(peer_upsert) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc); { char pn2[128]; json_flat_get(peer_userinfo, "name", pn2, sizeof(pn2)); topo_node_sqlite_node_update_verified(db, node_id, pn2, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } - chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); for (uint8_t i = 0; i < ac; i++) { diff --git a/tools/chatgui/transport/chat_sync.h b/tools/chatgui/transport/chat_sync.h index 03aa1622..0aab729f 100644 --- a/tools/chatgui/transport/chat_sync.h +++ b/tools/chatgui/transport/chat_sync.h @@ -54,11 +54,6 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, const uint8_t* addrs_data, int addr_count); -/* Авто-подключение к узлам группы (параллельно до 10, останавливается при 3 успешных) */ -void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst); -void chat_sync_auto_connect_stop(void); -void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id); - #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index dbbeb595..ea7ec9fe 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -28,6 +28,7 @@ extern "C" { #include "gui_bridge.h" #include "topo_node_sqlite.h" #include "topo_node.h" +#include "topo_group.h" } #define ETCP_RT_ID_CHAT 0x11 @@ -252,7 +253,6 @@ void UtunNode::runLoop() { } utun_instance_set_tun_init_enabled(0); - utun_instance_set_topo_group_enabled(0); struct UASYNC* ua = uasync_create(); if (!ua) { @@ -315,8 +315,27 @@ void UtunNode::runLoop() { /* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */ chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData()); chat_sync_init(m_instance, nullptr); - chat_sync_auto_connect_start(m_instance); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync + auto_connect initialized"); + + /* create TOPO_GROUP_TYPE_CHAT for all existing channels */ + if (m_instance->topo_groups) { + uint8_t buf[4096]; size_t ch_len = 0; + chat_core_list_channels(buf, sizeof(buf), &ch_len); + if (ch_len >= 2) { + uint16_t cnt; memcpy(&cnt, buf, 2); + const uint8_t* p = buf + 2; size_t rem = ch_len - 2; + for (uint16_t i = 0; i < cnt && rem >= 1; i++) { + uint8_t id_len = *p++; rem--; + if (rem < id_len) break; + char ch_id[64]; memcpy(ch_id, p, id_len < 63 ? id_len : 63); ch_id[id_len < 63 ? id_len : 63] = 0; + p += id_len; rem -= id_len; + uint64_t gid = strtoull(ch_id, NULL, 10); + if (gid != 0 && !topo_groups_find(m_instance->topo_groups, gid)) + topo_groups_create_group(m_instance->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id); + } + } + } + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop"); while (!m_stop) {