diff --git a/tools/chatgui/CMakeLists.txt b/tools/chatgui/CMakeLists.txt index 18a588cc..77fa56a5 100644 --- a/tools/chatgui/CMakeLists.txt +++ b/tools/chatgui/CMakeLists.txt @@ -78,6 +78,7 @@ 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 @@ -89,7 +90,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_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_conn_mgr.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/transport/chat_conn_mgr.c b/tools/chatgui/transport/chat_conn_mgr.c new file mode 100644 index 00000000..9b120e39 --- /dev/null +++ b/tools/chatgui/transport/chat_conn_mgr.c @@ -0,0 +1,589 @@ +/* + * chat_conn_mgr.c — менеджер ETCP-подключений (1 на инстанс) + * + * Управляет подключениями к узлам: перебор адресов, случайный выбор, + * ротация при падении соединения. До CM_MAX_PARALLEL одновременных + * подключений на узел. Фоновый пинг для измерения RTT. + */ + +#include "chat_conn_mgr.h" +#include "chat_core.h" +#include "gui_bridge.h" + +#include "../../../src/utun_instance.h" +#include "../../../src/etcp_api.h" +#include "../../../src/etcp.h" +#include "../../../src/etcp_connections.h" +#include "../../../src/secure_channel.h" +#include "../../../lib/u_async.h" +#include "../../../lib/ll_queue.h" +#include "../../../lib/mem.h" +#include "../../../lib/debug_config.h" + +#include +#include +#include +#include "../../../lib/platform_compat.h" + +#define CM_ID "chat_conn_mgr" + +/* ─── константы ─── */ +#define CM_MAX_PARALLEL 10 +#define CM_MAX_ADDRS 16 +#define CM_MAX_NODES 32 +#define CM_INIT_TIMEOUT_MS 3000 +#define CM_RETRY_MS 3000 +#define CM_PING_MS 5000 +#define CM_PING_TIMEOUT_MS 2000 +#define CM_INIT_TIMEOUT_TB (CM_INIT_TIMEOUT_MS * 10) +#define CM_RETRY_TB (CM_RETRY_MS * 10) +#define CM_PING_TB (CM_PING_MS * 10) + +/* ─── состояния слота ─── */ +#define CM_SLOT_FREE 0 +#define CM_SLOT_CONNECTING 1 +#define CM_SLOT_ACTIVE 2 +#define CM_SLOT_DEAD 3 + +/* ─── типы ─── */ + +struct cm_cb { + void (*cb)(int result, uint64_t node_id, void* arg); + void* arg; + struct cm_cb* next; +}; + +struct cm_slot { + uint8_t addr_idx; + uint8_t state; + uint16_t rtt; + struct ETCP_CONN* conn; + void* init_timer; +}; + +struct cm_addr { + struct sockaddr_storage sa; + uint16_t rtt; +}; + +struct cm_node; + +struct cm_slot_ctx { + struct cm_node* node; + int slot_idx; +}; + +struct cm_node { + uint64_t node_id; + uint8_t pubkey[SC_PUBKEY_SIZE]; + struct cm_addr addrs[CM_MAX_ADDRS]; + int addr_count; + int addr_loaded; + struct cm_slot slots[CM_MAX_PARALLEL]; + struct cm_cb* cbs; + int delivered; + void* retry_timer; + struct chat_conn_mgr* mgr; /* back-pointer */ +}; + +struct chat_conn_mgr { + struct UTUN_INSTANCE* inst; + struct sqlite3* db; + struct cm_node nodes[CM_MAX_NODES]; + int node_count; + void* ping_timer; +}; + +/* ─── forward ─── */ +static struct cm_node* cm_find_node(struct chat_conn_mgr* mgr, uint64_t node_id); +static struct cm_node* cm_ensure_node(struct chat_conn_mgr* mgr, uint64_t node_id); +static void cm_load_addrs(struct cm_node* node); +static int cm_count_active(struct cm_node* node); +static int cm_addr_taken(struct cm_node* node, int idx); +static int cm_external_conn_active(struct cm_node* node); +static struct cm_slot* cm_find_free(struct cm_node* node); +static struct cm_slot* cm_find_dead(struct cm_node* node); +static struct cm_slot* cm_find_slot_by_conn(struct cm_node* node, struct ETCP_CONN* conn); +static void cm_deliver_all(struct cm_node* node, int result); +static void cm_rotation(struct cm_node* node); +static void cm_init_cb(struct ETCP_CONN* conn, void* arg); +static void cm_init_timeout(void* arg); +static void cm_conn_down_cb(struct ETCP_CONN* conn, void* arg); +static void cm_retry_cb(void* arg); +static void cm_ping_timer_cb(void* arg); +static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len); +static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr); +static void cm_invite_result_cb(int result, uint64_t node_id, void* arg); + +/* ─── init / destroy ─── */ + +struct chat_conn_mgr* chat_conn_mgr_init(struct UTUN_INSTANCE* inst, struct sqlite3* db) { + if (!inst || !db) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: init bad args", CM_ID); return NULL; } + struct chat_conn_mgr* mgr = u_calloc(1, sizeof(struct chat_conn_mgr)); + if (!mgr) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: alloc failed", CM_ID); return NULL; } + mgr->inst = inst; mgr->db = db; + mgr->ping_timer = uasync_set_timeout(inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", CM_ID); + return mgr; +} + +void chat_conn_mgr_destroy(struct chat_conn_mgr* mgr) { + if (!mgr) return; + if (mgr->ping_timer) { uasync_cancel_timeout(mgr->inst->ua, mgr->ping_timer); mgr->ping_timer = NULL; } + for (int ni = 0; ni < mgr->node_count; ni++) { + struct cm_node* node = &mgr->nodes[ni]; + if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } + for (int si = 0; si < CM_MAX_PARALLEL; si++) { + struct cm_slot* slot = &node->slots[si]; + if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } + if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } + } + cm_deliver_all(node, CC_ERR_INTERNAL); + } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CM_ID); + u_free(mgr); +} + +struct sqlite3* chat_conn_mgr_get_db(struct chat_conn_mgr* mgr) { + return mgr ? mgr->db : NULL; +} + +/* ─── управление узлами ─── */ + +static struct cm_node* cm_find_node(struct chat_conn_mgr* mgr, uint64_t node_id) { + for (int i = 0; i < mgr->node_count; i++) + if (mgr->nodes[i].node_id == node_id) return &mgr->nodes[i]; + return NULL; +} + +static struct cm_node* cm_ensure_node(struct chat_conn_mgr* mgr, uint64_t node_id) { + struct cm_node* node = cm_find_node(mgr, node_id); + if (node) return node; + if (mgr->node_count >= CM_MAX_NODES) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: too many nodes", CM_ID); return NULL; } + node = &mgr->nodes[mgr->node_count++]; + memset(node, 0, sizeof(*node)); + node->node_id = node_id; node->mgr = mgr; + for (int i = 0; i < CM_MAX_PARALLEL; i++) node->slots[i].addr_idx = 255; + return node; +} + +void chat_conn_mgr_add_node(struct chat_conn_mgr* mgr, uint64_t node_id) { + if (!mgr) return; + if (cm_find_node(mgr, node_id)) return; + cm_ensure_node(mgr, node_id); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx added", CM_ID, (unsigned long long)node_id); +} + +/* ─── загрузка адресов ─── */ + +static void cm_load_addrs(struct cm_node* node) { + if (!node || node->addr_loaded) return; + struct sqlite3* db = node->mgr->db; + struct { uint8_t a[4]; uint16_t p; } raw[CM_MAX_ADDRS]; + int rc = 0; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, "SELECT address,port FROM node_addresses WHERE node_id=? AND family=4", + -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)node->node_id); + while (sqlite3_step(st) == SQLITE_ROW && rc < CM_MAX_ADDRS) { + const void* a = sqlite3_column_blob(st, 0); + if (a && sqlite3_column_bytes(st, 0) == 4) { memcpy(raw[rc].a, a, 4); raw[rc].p = (uint16_t)sqlite3_column_int(st, 1); rc++; } + } + sqlite3_finalize(st); + } + int uniq = 0; + for (int i = 0; i < rc; i++) { + int dup = 0; + for (int j = 0; j < uniq; j++) if (memcmp(raw[i].a, raw[j].a, 4) == 0 && raw[i].p == raw[j].p) { dup = 1; break; } + if (!dup) { if (i != uniq) raw[uniq] = raw[i]; uniq++; } + } + for (int i = 0; i < uniq; i++) { + struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, raw[i].a, 4); sin.sin_port = htons(raw[i].p); + memcpy(&node->addrs[i].sa, &sin, sizeof(sin)); node->addrs[i].rtt = 65535; + } + node->addr_count = uniq; node->addr_loaded = 1; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx loaded %d addrs", CM_ID, (unsigned long long)node->node_id, uniq); + + uint8_t pubkey[SC_PUBKEY_SIZE] = {0}; + if (sqlite3_prepare_v2(db, "SELECT x25519_pubkey FROM nodes WHERE node_id=?", + -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)node->node_id); + if (sqlite3_step(st) == SQLITE_ROW) { + const void* pk = sqlite3_column_blob(st, 0); + if (pk && sqlite3_column_bytes(st, 0) == SC_PUBKEY_SIZE) memcpy(pubkey, pk, SC_PUBKEY_SIZE); + } + sqlite3_finalize(st); + } + if (pubkey[0] || pubkey[1]) memcpy(node->pubkey, pubkey, SC_PUBKEY_SIZE); + else DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx no pubkey in DB", CM_ID, (unsigned long long)node->node_id); +} + +/* ─── helpers ─── */ + +static int cm_count_active(struct cm_node* node) { + int a = 0; + for (int i = 0; i < CM_MAX_PARALLEL; i++) if (node->slots[i].state == CM_SLOT_ACTIVE) a++; + return a; +} + +static int cm_addr_taken(struct cm_node* node, int idx) { + for (int i = 0; i < CM_MAX_PARALLEL; i++) + if (node->slots[i].state != CM_SLOT_FREE && node->slots[i].addr_idx == idx) return 1; + return 0; +} + +static int cm_external_conn_active(struct cm_node* node) { + struct UTUN_INSTANCE* inst = node->mgr->inst; + struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node->node_id); + if (!e) return 0; + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + return (ce && ce->conn && ce->conn->links_up && ce->conn->peer_node_id == node->node_id); +} + +static struct cm_slot* cm_find_free(struct cm_node* node) { + for (int i = 0; i < CM_MAX_PARALLEL; i++) + if (node->slots[i].state == CM_SLOT_FREE) return &node->slots[i]; + return NULL; +} + +static struct cm_slot* cm_find_dead(struct cm_node* node) { + for (int i = 0; i < CM_MAX_PARALLEL; i++) + if (node->slots[i].state == CM_SLOT_DEAD) return &node->slots[i]; + return NULL; +} + +static struct cm_slot* cm_find_slot_by_conn(struct cm_node* node, struct ETCP_CONN* conn) { + for (int i = 0; i < CM_MAX_PARALLEL; i++) + if (node->slots[i].conn == conn) return &node->slots[i]; + return NULL; +} + +static void cm_deliver_all(struct cm_node* node, int result) { + while (node->cbs) { struct cm_cb* cb = node->cbs; node->cbs = cb->next; cb->cb(result, node->node_id, cb->arg); u_free(cb); } +} + +static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr) { + struct ETCP_SOCKET* s = mgr->inst->etcp_sockets; + while (s) { if (s->local_addr.ss_family == AF_INET) return s; s = s->next; } + return NULL; +} + +/* ─── ротация ─── */ + +static void cm_rotation(struct cm_node* node) { + if (!node) return; + struct chat_conn_mgr* mgr = node->mgr; + + int active = cm_count_active(node); + + if (cm_external_conn_active(node)) { + if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } + if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } + return; + } + + if (active >= CM_MAX_PARALLEL) { + if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } + if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } + return; + } + + struct cm_slot* slot = cm_find_free(node); + if (!slot) slot = cm_find_dead(node); + if (!slot) return; + + struct ETCP_SOCKET* sock = cm_best_socket(mgr); + if (!sock) return; + + for (int attempt = 0; attempt < 3; attempt++) { + int idx = rand() % node->addr_count; + if (!cm_addr_taken(node, idx)) { + slot->addr_idx = idx; slot->state = CM_SLOT_CONNECTING; + struct ETCP_CONN* conn = etcp_connection_create(mgr->inst, NULL); + if (!conn) { slot->state = CM_SLOT_FREE; return; } + sc_init_ctx(&conn->crypto_ctx, &mgr->inst->my_keys); + sc_set_peer_public_key(&conn->crypto_ctx, node->pubkey, 0); + etcp_conn_set_peer_node_id(conn, node->node_id); + struct cm_slot_ctx* ctx = u_calloc(1, sizeof(struct cm_slot_ctx)); + if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } + ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); + etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); + if (!etcp_link_new(conn, sock, &node->addrs[idx].sa, 0)) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, idx); + etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; + cm_rotation(node); return; + } + slot->conn = conn; + slot->init_timer = uasync_set_timeout(mgr->inst->ua, CM_INIT_TIMEOUT_TB, ctx, cm_init_timeout, "cm_init"); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT node=0x%016llx slot=%d addr=%d", CM_ID, (unsigned long long)node->node_id, ctx->slot_idx, idx); + return; + } + } + + if (slot->addr_idx != 255) { + slot->state = CM_SLOT_CONNECTING; + struct ETCP_CONN* conn = etcp_connection_create(mgr->inst, NULL); + if (!conn) { slot->state = CM_SLOT_FREE; return; } + sc_init_ctx(&conn->crypto_ctx, &mgr->inst->my_keys); + sc_set_peer_public_key(&conn->crypto_ctx, node->pubkey, 0); + etcp_conn_set_peer_node_id(conn, node->node_id); + struct cm_slot_ctx* ctx = u_calloc(1, sizeof(struct cm_slot_ctx)); + if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } + ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); + etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); + if (!etcp_link_new(conn, sock, &node->addrs[slot->addr_idx].sa, 0)) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new(retry) failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, slot->addr_idx); + etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; + cm_rotation(node); return; + } + slot->conn = conn; + slot->init_timer = uasync_set_timeout(mgr->inst->ua, CM_INIT_TIMEOUT_TB, ctx, cm_init_timeout, "cm_init"); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT(retry) node=0x%016llx slot=%d addr=%d", CM_ID, (unsigned long long)node->node_id, ctx->slot_idx, slot->addr_idx); + return; + } + + slot->state = CM_SLOT_FREE; +} + +/* ─── коллбэки ─── */ + +static void cm_init_cb(struct ETCP_CONN* conn, void* arg) { + struct cm_slot_ctx* ctx = (struct cm_slot_ctx*)arg; + struct cm_node* node = ctx->node; + struct cm_slot* slot = &node->slots[ctx->slot_idx]; + + if (slot->init_timer) { uasync_cancel_timeout(node->mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } + + if (slot->conn && slot->conn != conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } + + slot->state = CM_SLOT_ACTIVE; slot->conn = conn; slot->rtt = conn->rtt_last; + etcp_conn_add_down_cbk(conn, cm_conn_down_cb, node); + etcp_conn_remove_init_cbk(conn, cm_init_cb, ctx); + u_free(ctx); + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT OK node=0x%016llx slot=%d rtt=%u active=%d", + CM_ID, (unsigned long long)node->node_id, (int)(slot - node->slots), slot->rtt, cm_count_active(node)); + + struct chat_conn_mgr* mgr = node->mgr; + int active = cm_count_active(node); + + if (active > CM_MAX_PARALLEL) { + int worst_idx = -1; uint16_t worst_rtt = 0; + for (int i = 0; i < CM_MAX_PARALLEL; i++) { + if (node->slots[i].state == CM_SLOT_ACTIVE && node->slots[i].rtt > worst_rtt) { + worst_rtt = node->slots[i].rtt; worst_idx = i; + } + } + if (worst_idx >= 0) { + struct cm_slot* w = &node->slots[worst_idx]; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: evict slot=%d rtt=%u node=0x%016llx", + CM_ID, worst_idx, w->rtt, (unsigned long long)node->node_id); + w->state = CM_SLOT_FREE; w->conn = NULL; w->addr_idx = 255; + } + active = cm_count_active(node); + } + + if (active >= CM_MAX_PARALLEL) { + if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } + if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } + return; + } + + cm_rotation(node); +} + +static void cm_init_timeout(void* arg) { + struct cm_slot_ctx* ctx = (struct cm_slot_ctx*)arg; + struct cm_node* node = ctx->node; + struct cm_slot* slot = &node->slots[ctx->slot_idx]; + slot->init_timer = NULL; + + if (slot->state == CM_SLOT_CONNECTING) slot->state = CM_SLOT_DEAD; + node->delivered = 0; + + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT TIMEOUT node=0x%016llx slot=%d", + CM_ID, (unsigned long long)node->node_id, ctx->slot_idx); + + cm_rotation(node); + u_free(ctx); +} + +static void cm_conn_down_cb(struct ETCP_CONN* conn, void* arg) { + struct cm_node* node = (struct cm_node*)arg; + struct cm_slot* slot = cm_find_slot_by_conn(node, conn); + if (!slot) return; + + slot->state = CM_SLOT_FREE; slot->conn = NULL; + node->delivered = 0; + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: DOWN node=0x%016llx slot=%d active=%d", + CM_ID, (unsigned long long)node->node_id, (int)(slot - node->slots), cm_count_active(node)); + + cm_rotation(node); +} + +static void cm_retry_cb(void* arg) { + struct cm_node* node = (struct cm_node*)arg; + node->retry_timer = NULL; + + for (int i = 0; i < CM_MAX_PARALLEL; i++) + if (node->slots[i].state == CM_SLOT_DEAD) { node->slots[i].state = CM_SLOT_FREE; node->slots[i].addr_idx = 255; } + + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry node=0x%016llx active=%d", CM_ID, (unsigned long long)node->node_id, cm_count_active(node)); + + cm_rotation(node); + + if (cm_count_active(node) < CM_MAX_PARALLEL && !node->delivered) + node->retry_timer = uasync_set_timeout(node->mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry"); +} + +/* ─── фоновый пинг ─── */ + +static void cm_ping_timer_cb(void* arg) { + struct chat_conn_mgr* mgr = (struct chat_conn_mgr*)arg; + if (mgr->node_count == 0) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } + + int ni = rand() % mgr->node_count; + struct cm_node* node = &mgr->nodes[ni]; + if (!node->addr_loaded || node->addr_count == 0) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } + + int ai = rand() % node->addr_count; + + if (cm_addr_taken(node, ai)) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } + + struct ETCP_SOCKET* sock = cm_best_socket(mgr); + if (sock) etcp_send_ping_to_socket(mgr->inst, sock, node->pubkey, &node->addrs[ai].sa, CM_PING_TIMEOUT_MS, cm_ping_result_cb, &node->addrs[ai], NULL, 0); + + mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); +} + +static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len) { + (void)nonce; (void)resp_data; (void)resp_data_len; + struct cm_addr* addr = (struct cm_addr*)arg; + addr->rtt = success ? rtt : 65535; +} + +/* ─── публичное API ─── */ + +void chat_conn_mgr_connect(struct chat_conn_mgr* mgr, uint64_t node_id, + void (*cb)(int result, uint64_t node_id, void* arg), + void* arg) { + if (!mgr) { if (cb) cb(CC_ERR_INTERNAL, node_id, arg); return; } + + struct cm_node* node = cm_ensure_node(mgr, node_id); + if (!node) { if (cb) cb(CC_ERR_INTERNAL, node_id, arg); return; } + + if (!node->addr_loaded) cm_load_addrs(node); + + if (cb) { + struct cm_cb* cb_node = u_calloc(1, sizeof(struct cm_cb)); + if (cb_node) { cb_node->cb = cb; cb_node->arg = arg; cb_node->next = node->cbs; node->cbs = cb_node; } + } + + if (node->pubkey[0] == 0 && node->pubkey[1] == 0) { + DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect to 0x%016llx — no pubkey in DB", CM_ID, (unsigned long long)node_id); + if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_ERR_NOT_FOUND); } + return; + } + if (node->addr_count == 0) { + DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect to 0x%016llx — no addresses", CM_ID, (unsigned long long)node_id); + if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_ERR_NO_ADDRESSES); } + return; + } + + if (node->delivered) return; + + if (cm_external_conn_active(node)) { + node->delivered = 1; cm_deliver_all(node, CC_OK); + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx already connected externally", CM_ID, (unsigned long long)node_id); + return; + } + + if (cm_count_active(node) >= 1) { + // at least one active already (might be from a previous rotation that didn't deliver yet) + if (cm_count_active(node) >= CM_MAX_PARALLEL && !node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } + return; + } + + cm_rotation(node); + + if (cm_count_active(node) == 0 && !node->retry_timer && !node->delivered) + node->retry_timer = uasync_set_timeout(mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry"); +} + +void chat_conn_mgr_cancel(struct chat_conn_mgr* mgr, uint64_t node_id) { + if (!mgr) return; + struct cm_node* node = cm_find_node(mgr, node_id); + if (!node) return; + + if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } + for (int i = 0; i < CM_MAX_PARALLEL; i++) { + struct cm_slot* slot = &node->slots[i]; + if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } + if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } + slot->state = CM_SLOT_FREE; slot->addr_idx = 255; + } + cm_deliver_all(node, CC_ERR_INTERNAL); + node->delivered = 0; + DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: cancelled node 0x%016llx", CM_ID, (unsigned long long)node_id); +} + +void chat_conn_mgr_connect_from_invite(struct chat_conn_mgr* mgr, struct chat_invite* inv) { + if (!mgr || !inv) return; + uint64_t node_id = inv->node_id; + + /* save pubkey */ + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(mgr->db, "INSERT INTO nodes(node_id,x25519_pubkey) VALUES(?,?)" + " ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey", + -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); + sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); + sqlite3_step(st); sqlite3_finalize(st); + } + + /* save addresses (socket_id=0) */ + sqlite3_stmt* ds = NULL; + if (sqlite3_prepare_v2(mgr->db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", + -1, &ds, NULL) == SQLITE_OK) { + sqlite3_bind_int64(ds, 1, (sqlite3_int64)node_id); + sqlite3_step(ds); sqlite3_finalize(ds); + } + sqlite3_stmt* is = NULL; + if (sqlite3_prepare_v2(mgr->db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id)" + " VALUES(?,?,1,?,?,0,?)", -1, &is, NULL) == SQLITE_OK) { + const uint8_t* src = inv->addrs_data; + for (int i = 0; i < inv->addr_count; i++) { + uint8_t family = *src++; uint8_t sid = *src++; + if (family == 4) { + sqlite3_bind_int64(is, 1, (sqlite3_int64)node_id); + sqlite3_bind_int(is, 2, 4); + sqlite3_bind_blob(is, 3, src, 4, SQLITE_STATIC); src += 4; + uint16_t port = ((uint16_t)src[0] << 8) | src[1]; src += 2; + sqlite3_bind_int(is, 4, (int)port); + sqlite3_bind_int(is, 5, (int)sid); + sqlite3_step(is); sqlite3_reset(is); + } else src += 18; + } + sqlite3_finalize(is); + } + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: invite node=0x%016llx addrs=%d — saved, connecting", + CM_ID, (unsigned long long)node_id, inv->addr_count); + + /* add node + reset addr_loaded to reload from DB */ + struct cm_node* node = cm_ensure_node(mgr, node_id); + if (node) { node->addr_loaded = 0; cm_load_addrs(node); } + + { char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)inv->channel_id); + chat_core_ensure_channel_ready(ch_str); } + + chat_conn_mgr_connect(mgr, node_id, cm_invite_result_cb, &inv->channel_id); +} + +static void cm_invite_result_cb(int result, uint64_t node_id, void* arg) { + uint64_t channel_id = arg ? *(uint64_t*)arg : 0; + uint8_t data[20]; memcpy(data, &node_id, 8); memcpy(data + 8, &result, 4); memcpy(data + 12, &channel_id, 8); + gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); +} diff --git a/tools/chatgui/transport/chat_conn_mgr.h b/tools/chatgui/transport/chat_conn_mgr.h new file mode 100644 index 00000000..5539807a --- /dev/null +++ b/tools/chatgui/transport/chat_conn_mgr.h @@ -0,0 +1,55 @@ +/* + * chat_conn_mgr.h — менеджер ETCP-подключений (1 на инстанс) + * + * Управляет подключениями к узлам: перебор адресов, случайный выбор, + * ротация при падении соединения. До CM_MAX_PARALLEL одновременных + * подключений на узел. Фоновый пинг для измерения RTT. + */ + +#ifndef CHAT_CONN_MGR_H +#define CHAT_CONN_MGR_H + +#include +#include + +struct sqlite3; +struct UTUN_INSTANCE; + +/* ── Коды возврата ── */ +#define CC_OK 0 +#define CC_ERR_NOT_FOUND -1 +#define CC_ERR_NO_ADDRESSES -2 +#define CC_ERR_TIMEOUT -3 +#define CC_ERR_UNREACHABLE -4 +#define CC_ERR_INTERNAL -7 + +/* ── Жизненный цикл ── */ + +struct chat_conn_mgr* chat_conn_mgr_init(struct UTUN_INSTANCE* inst, struct sqlite3* db); +void chat_conn_mgr_destroy(struct chat_conn_mgr* mgr); +struct sqlite3* chat_conn_mgr_get_db(struct chat_conn_mgr* mgr); + +/* ── Управление узлами ── */ + +void chat_conn_mgr_add_node(struct chat_conn_mgr* mgr, uint64_t node_id); + +/* ── Подключение ── */ + +void chat_conn_mgr_connect(struct chat_conn_mgr* mgr, uint64_t node_id, + void (*cb)(int result, uint64_t node_id, void* arg), + void* arg); +void chat_conn_mgr_cancel(struct chat_conn_mgr* mgr, uint64_t node_id); + +/* ── invite ── */ + +struct chat_invite { + uint64_t channel_id; + uint64_t node_id; + uint8_t pubkey[32]; + uint8_t* addrs_data; + int addr_count; +}; + +void chat_conn_mgr_connect_from_invite(struct chat_conn_mgr* mgr, struct chat_invite* inv); + +#endif /* CHAT_CONN_MGR_H */ diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 708e1642..95e27863 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -6,6 +6,7 @@ */ #include "chat_core.h" +#include "chat_conn_mgr.h" #include "db_sync.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" @@ -39,6 +40,7 @@ static struct chat_core_ctx { sqlite3* db; uint8_t shared_db; /* db == inst->topo_sqlite_db, do not close */ uint64_t my_node_id; + struct chat_conn_mgr* conn_mgr; /* db_sync instances per channel */ struct DB_SYNC_INSTANCE** si; @@ -139,6 +141,7 @@ 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(); @@ -230,14 +233,26 @@ sqlite3* chat_core_get_db(void) { return g_cc.db; } +struct UTUN_INSTANCE* chat_core_get_inst(void) { + return g_cc.inst; +} + +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; + g_cc.db = NULL; g_cc.inst = NULL; DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID); } @@ -625,327 +640,6 @@ int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, return 0; } -/* ─── подключение к пиру из invite-ссылки ─── */ - -/* проверить, есть ли уже ETCP-подключение к узлу (готовое или pending) */ -static int chat_core_has_conn(uint64_t node_id) { - if (!g_cc.inst || !g_cc.inst->connections) return 0; - struct ll_entry* entry = g_cc.inst->connections->head; - while (entry) { - struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; - if (ce->conn->peer_node_id == node_id) return 1; - entry = entry->next; - } - return 0; -} - -static void connect_result_cb(int result, uint64_t node_id, void* arg) { - uint64_t channel_id = arg ? *(uint64_t*)arg : 0; - uint8_t data[20]; - memcpy(data, &node_id, 8); - memcpy(data + 8, &result, 4); - memcpy(data + 12, &channel_id, 8); - gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); -} - -void chat_core_connect_from_invite(struct chat_invite* inv) { - if (!g_cc.initialized || !g_cc.inst || !inv) return; - - uint64_t node_id = inv->node_id; - - /* save pubkey to nodes table */ - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(g_cc.db, - "INSERT INTO nodes(node_id, x25519_pubkey) VALUES(?,?) ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); - sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); - sqlite3_step(st); sqlite3_finalize(st); - } - - /* save addresses to node_addresses */ - { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] save_invite_addr node=0x%016llx addr_count=%d", CC_ID, (unsigned long long)node_id, inv->addr_count); - sqlite3_stmt* ds = NULL; - if (sqlite3_prepare_v2(g_cc.db, - "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", -1, &ds, NULL) == SQLITE_OK) { - sqlite3_bind_int64(ds, 1, (sqlite3_int64)node_id); - sqlite3_step(ds); - int deleted = sqlite3_changes(g_cc.db); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] invite DELETE socket_id=0 for node=0x%016llx: %d rows deleted", CC_ID, (unsigned long long)node_id, deleted); - sqlite3_finalize(ds); - } - sqlite3_stmt* is = NULL; - if (sqlite3_prepare_v2(g_cc.db, - "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id)" - " VALUES(?,?,1,?,?,0,?)", -1, &is, NULL) == SQLITE_OK) { - const uint8_t* src = inv->addrs_data; - int written = 0; - for (int i = 0; i < inv->addr_count; i++) { - uint8_t family = *src++; - uint8_t sid = *src++; - if (family == 4) { - sqlite3_bind_int64(is, 1, (sqlite3_int64)node_id); - sqlite3_bind_int(is, 2, 4); - sqlite3_bind_blob(is, 3, src, 4, SQLITE_STATIC); src += 4; - uint16_t port = ((uint16_t)src[0] << 8) | src[1]; src += 2; - sqlite3_bind_int(is, 4, (int)port); - sqlite3_bind_int(is, 5, (int)sid); - sqlite3_step(is); sqlite3_reset(is); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] invite INSERT node=0x%016llx sock=%d %d.%d.%d.%d:%d", - CC_ID, (unsigned long long)node_id, sid, src[-6], src[-5], src[-4], src[-3], port); - written++; - } else { - src += 18; - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] invite SKIP v6 addr for node=0x%016llx", CC_ID, (unsigned long long)node_id); - } - } - sqlite3_finalize(is); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] invite DONE: %d v4 addresses written for node=0x%016llx", - CC_ID, written, (unsigned long long)node_id); - } else { - DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] invite FAIL prepare INSERT node=0x%016llx", CC_ID, (unsigned long long)node_id); - } - } - - if (chat_core_has_conn(node_id)) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: invite node 0x%016llx already connected, reusing", - CC_ID, (unsigned long long)node_id); - connect_result_cb(CC_OK, node_id, &inv->channel_id); - return; - } - - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: invite saved node=0x%016llx addrs=%d, starting direct connect", - CC_ID, (unsigned long long)node_id, inv->addr_count); - - { char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)inv->channel_id); - chat_core_ensure_channel_ready(ch_str); } - - chat_core_connect_auto(node_id, connect_result_cb, &inv->channel_id, NULL); -} - -/* ═══ Auto-connect: direct ETCP connection from SQLite (no BGP) ═══ */ - -#define CA_CONNECT_TIMEOUT_MS 3000 - -struct ca_state; - -struct ca_ctx { - struct ca_state* state; - int addr_index; -}; - -struct ca_state { - struct UTUN_INSTANCE* inst; - int addr_count; - int pending_count; - int delivered; /* 0=pending, 1=result already delivered */ - int cancelled; /* 1=externally cancelled, do not deliver result */ - uint64_t node_id; - struct ETCP_CONN** conns; - void** timers; - struct ca_ctx** ctxs; - void (*result_cb)(int result, uint64_t node_id, void* arg); - void* result_arg; - uint8_t pubkey[SC_PUBKEY_SIZE]; -}; - -static void ca_cleanup(struct ca_state* st) { - if (!st) return; - if (st->ctxs) { for (int i = 0; i < st->addr_count; i++) u_free(st->ctxs[i]); u_free(st->ctxs); } - u_free(st->conns); - u_free(st->timers); - u_free(st); -} - -static void ca_init_cb(struct ETCP_CONN* conn, void* arg) { - struct ca_ctx* ctx = (struct ca_ctx*)arg; - struct ca_state* st = ctx->state; - if (st->delivered || st->cancelled) return; - st->delivered = 1; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect SUCCESS idx=%d peer=0x%016llx", - CC_ID, ctx->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 != ctx->addr_index) { uasync_call_soon(st->inst->ua, st->conns[i], (timeout_callback_t)etcp_connection_close); st->conns[i] = NULL; } - } - etcp_conn_remove_init_cbk(conn, ca_init_cb, ctx); - st->result_cb(CC_OK, st->node_id, st->result_arg); - ca_cleanup(st); -} - -static void ca_timeout_cb(void* arg) { - struct ca_ctx* ctx = (struct ca_ctx*)arg; - struct ca_state* st = ctx->state; - if (st->delivered || st->cancelled) return; - st->timers[ctx->addr_index] = NULL; - st->pending_count--; - DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx", - CC_ID, ctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id); - if (st->conns[ctx->addr_index]) { - etcp_connection_close(st->conns[ctx->addr_index]); - st->conns[ctx->addr_index] = NULL; - } - if (st->pending_count <= 0 && !st->delivered) { - st->delivered = 1; - 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; } - } - st->result_cb(CC_ERR_TIMEOUT, st->node_id, st->result_arg); - ca_cleanup(st); - } -} - -void chat_core_connect_auto_cancel(void* state) { - if (!state) return; - struct ca_state* st = (struct ca_state*)state; - st->cancelled = 1; - DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect CANCELLED peer=0x%016llx", - CC_ID, (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]) { etcp_connection_close(st->conns[i]); st->conns[i] = NULL; } - } - ca_cleanup(st); -} - -void chat_core_connect_auto(uint64_t node_id, - void (*cb)(int result, uint64_t node_id, void* arg), - void* arg, - void** out_state) { - if (!g_cc.initialized || !g_cc.inst || !cb) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (not initialized) node=0x%016llx init=%d inst=%p cb=%p", - CC_ID, (unsigned long long)node_id, g_cc.initialized, (void*)g_cc.inst, (void*)cb); - return; - } - if (out_state) *out_state = NULL; - - 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_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (no AF_INET socket) node=0x%016llx sockets=%p", - CC_ID, (unsigned long long)node_id, (void*)g_cc.inst->etcp_sockets); - cb(CC_ERR_INTERNAL, node_id, arg); - return; - } - - uint8_t pubkey[SC_PUBKEY_SIZE] = {0}; - sqlite3_stmt* st = NULL; - if (sqlite3_prepare_v2(g_cc.db, - "SELECT x25519_pubkey FROM nodes WHERE node_id=?", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); - if (sqlite3_step(st) == SQLITE_ROW) { - const void* pk = sqlite3_column_blob(st, 0); - if (pk && sqlite3_column_bytes(st, 0) == SC_PUBKEY_SIZE) - memcpy(pubkey, pk, SC_PUBKEY_SIZE); - } - sqlite3_finalize(st); - } - if (pubkey[0] == 0 && pubkey[1] == 0) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (no pubkey in DB) node=0x%016llx", - CC_ID, (unsigned long long)node_id); - cb(CC_ERR_NOT_FOUND, node_id, arg); - return; - } - - /* collect IPv4 addresses */ - struct { uint8_t addr[4]; uint16_t port; } addrs[16]; - int addr_count = 0; - if (sqlite3_prepare_v2(g_cc.db, - "SELECT address, port FROM node_addresses WHERE node_id=? AND family=4", - -1, &st, NULL) == SQLITE_OK) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); - while (sqlite3_step(st) == SQLITE_ROW && addr_count < 16) { - const void* a = sqlite3_column_blob(st, 0); - int alen = sqlite3_column_bytes(st, 0); - if (a && alen == 4) { - memcpy(addrs[addr_count].addr, a, 4); - addrs[addr_count].port = (uint16_t)sqlite3_column_int(st, 1); - addr_count++; - } - } - sqlite3_finalize(st); - } - /* dedup by (addr, port) — can't have two links to same socket */ - { int uniq = 0; - for (int i = 0; i < addr_count; i++) { - int dup = 0; - for (int j = 0; j < uniq; j++) - if (memcmp(addrs[i].addr, addrs[j].addr, 4) == 0 && addrs[i].port == addrs[j].port) { dup = 1; break; } - if (!dup) { if (i != uniq) addrs[uniq] = addrs[i]; uniq++; } - } - addr_count = uniq; } - if (addr_count == 0) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (no IPv4 addrs in DB) node=0x%016llx (BGP not synced?)", - CC_ID, (unsigned long long)node_id); - cb(CC_ERR_NO_ADDRESSES, node_id, arg); - return; - } - - struct ca_state* pst = u_calloc(1, sizeof(struct ca_state)); - if (!pst) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (alloc ca_state) node=0x%016llx", CC_ID, (unsigned long long)node_id); cb(CC_ERR_INTERNAL, node_id, arg); return; } - pst->inst = g_cc.inst; - pst->addr_count = addr_count; - pst->pending_count = addr_count; - pst->node_id = node_id; - pst->result_cb = cb; - pst->result_arg = arg; - memcpy(pst->pubkey, pubkey, SC_PUBKEY_SIZE); - pst->conns = u_calloc(addr_count, sizeof(struct ETCP_CONN*)); - pst->timers = u_calloc(addr_count, sizeof(void*)); - pst->ctxs = u_calloc(addr_count, sizeof(struct ca_ctx*)); - if (!pst->conns || !pst->timers || !pst->ctxs) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (alloc arrays) node=0x%016llx", CC_ID, (unsigned long long)node_id); ca_cleanup(pst); cb(CC_ERR_INTERNAL, node_id, arg); return; } - - for (int i = 0; i < addr_count; i++) { - struct sockaddr_in sin; - memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; - memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); - sin.sin_port = htons(addrs[i].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_CONNECTIVITY, "%s: auto_connect etcp_connection_create failed idx=%d", CC_ID, i); - pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; - } - sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); - sc_set_peer_public_key(&conn->crypto_ctx, pubkey, 0); - etcp_conn_set_peer_node_id(conn, node_id); - - struct ca_ctx* pctx = u_calloc(1, sizeof(struct ca_ctx)); - if (!pctx) { etcp_connection_close(conn); pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; } - pctx->state = pst; pctx->addr_index = i; - pst->ctxs[i] = pctx; - - etcp_conn_add_init_cbk(conn, ca_init_cb, pctx); - - if (!etcp_link_new(conn, best_socket, &sa, 0)) { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect etcp_link_new failed idx=%d", CC_ID, i); - u_free(pctx); etcp_connection_close(conn); - pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; - } - pst->conns[i] = conn; - pst->timers[i] = uasync_set_timeout(g_cc.inst->ua, CA_CONNECT_TIMEOUT_MS * 10, - pctx, ca_timeout_cb, "ca_timeout"); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect attempt %d/%d to %d.%d.%d.%d:%d", - CC_ID, i + 1, addr_count, - addrs[i].addr[0], addrs[i].addr[1], addrs[i].addr[2], addrs[i].addr[3], - addrs[i].port); - } - - if (out_state) *out_state = pst; - - if (pst->pending_count <= 0 && !pst->delivered) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: auto_connect FAIL (all %d attempts failed to start) node=0x%016llx", - CC_ID, addr_count, (unsigned long long)node_id); - if (out_state) *out_state = NULL; - cb(CC_ERR_UNREACHABLE, node_id, arg); - ca_cleanup(pst); - } -} - /* ─── подготовка инфраструктуры канала (db_sync instance) ─── */ void chat_core_ensure_channel_ready(const char* ch_id) { diff --git a/tools/chatgui/transport/chat_core.h b/tools/chatgui/transport/chat_core.h index 3eda1ba4..3dbb06ca 100644 --- a/tools/chatgui/transport/chat_core.h +++ b/tools/chatgui/transport/chat_core.h @@ -14,20 +14,16 @@ struct sqlite3; struct UTUN_INSTANCE; -/* ── Коды возврата подключения ── */ -#define CC_OK 0 -#define CC_ERR_NOT_FOUND -1 -#define CC_ERR_NO_ADDRESSES -2 -#define CC_ERR_TIMEOUT -3 -#define CC_ERR_UNREACHABLE -4 -#define CC_ERR_INTERNAL -7 - /* ── Жизненный цикл ── */ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path); void chat_core_destroy(struct UTUN_INSTANCE* inst); 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) ── */ @@ -49,30 +45,6 @@ int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len) int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len); -/* ── Подключение к пиру (из invite-ссылки) ── */ - -struct chat_invite { - uint64_t channel_id; - uint64_t node_id; - uint8_t pubkey[32]; - uint8_t* addrs_data; - int addr_count; -}; - -void chat_core_connect_from_invite(struct chat_invite* inv); - -/* Прямое ETCP-подключение к узлу без BGP: pubkey+адреса из SQLite, - * коллбэк вызывается один раз с CONN_MGR_OK / CONN_MGR_ERR_*. - * При out_state != NULL — возвращает opaque handle для отмены. */ -void chat_core_connect_auto(uint64_t node_id, - void (*cb)(int result, uint64_t node_id, void* arg), - void* arg, - void** out_state); - -/* Отмена активного auto-connect (закрывает ETCP-соединения, освобождает память). - * Коллбэк chat_core_connect_auto после cancel не вызывается. */ -void chat_core_connect_auto_cancel(void* state); - /* ── Создание канала (GUI → uasync) ── */ struct chat_channel_create { diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index 6c4b26a5..b0abc785 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -1,5 +1,6 @@ #include "chat_sync.h" #include "chat_core.h" +#include "chat_conn_mgr.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" #include "member_sync.h" @@ -33,8 +34,8 @@ static struct chat_sync* g_cs = NULL; struct ac_flight { uint64_t node_id; - void* ca_state; /* opaque, owned by chat_core_connect_auto */ - uint64_t created_tb; /* get_time_tb() when launched */ + uint8_t active; /* 1 = connection attempt in progress */ + uint64_t created_tb; }; struct auto_connect { @@ -86,7 +87,7 @@ 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].ca_state) continue; + 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) */ @@ -104,13 +105,13 @@ static void ac_gc(struct auto_connect* ac) { } } if (found_up) { - chat_core_connect_auto_cancel(ac->flights[i].ca_state); + 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_core_connect_auto_cancel(ac->flights[i].ca_state); + chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id); memset(&ac->flights[i], 0, sizeof(ac->flights[i])); } } @@ -119,14 +120,14 @@ static void ac_gc(struct auto_connect* ac) { 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].ca_state) n++; + 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].ca_state) return i; + for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (!ac->flights[i].active) return i; return -1; } @@ -172,16 +173,16 @@ static void ac_fill(struct auto_connect* ac) { /* already in a flight slot? */ int dup = 0; for (int i = 0; i < AC_MAX_FLIGHTS; i++) - if (ac->flights[i].ca_state && ac->flights[i].node_id == nid) { dup = 1; break; } + 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_core_connect_auto(nid, ac_result_cb, f, &f->ca_state); - if (!f->ca_state) { - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: flight to 0x%016llx FAIL (ca_state=NULL, sync fail in chat_core_connect_auto)", AC_ID, (unsigned long long)nid); + 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; } @@ -192,14 +193,14 @@ static void ac_fill(struct auto_connect* ac) { } } -/* ── result callback (fired by chat_core_connect_auto) ── */ +/* ── 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->ca_state = NULL; /* already freed by ca_init_cb or ca_timeout_cb */ + f->active = 0; f->node_id = 0; f->created_tb = 0; } /* ── retry timer callback (every AC_RETRY_MS) ── */ @@ -239,7 +240,7 @@ void chat_sync_auto_connect_stop(void) { 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].ca_state) { chat_core_connect_auto_cancel(ac->flights[i].ca_state); memset(&ac->flights[i], 0, sizeof(ac->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); @@ -802,7 +803,15 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { (void)inst; - chat_core_connect_auto(node_id, NULL, NULL, NULL); + chat_conn_mgr_connect(chat_conn_mgr_get(), node_id, NULL, NULL); +} + +struct cm_invite_wrap { struct chat_conn_mgr* mgr; 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); + u_free(w); } void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, @@ -836,8 +845,10 @@ 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); gui_bridge_post_uasync_fn( - (void(*)(void*))chat_core_connect_from_invite, inv); + (void(*)(void*))cm_invite_trampoline, w); } /* ─── Ed25519 sign / verify helpers ─── */ @@ -1191,6 +1202,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer, int rc = topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, join_ts, NULL, 0, x25519, ed_pub, joiner_name, NULL); 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); topo_node_sqlite_node_update_verified(db, node_id, joiner_name, x25519, ed_pub, join_ts); + chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id); /* save joiner node_info to local DB */ if (db && joiner_name[0]) { sqlite3_stmt* ns = NULL; @@ -1357,6 +1369,8 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(welcome) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc); topo_node_sqlite_node_update_verified(db, node_id, peer_name, x25519, ed_pub, join_ts); + 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++) { if (p + 2 > pl + len) break; @@ -1452,6 +1466,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, int rc = topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, join_ts, NULL, 0, x25519, ed_pub, peer_name, NULL); 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); topo_node_sqlite_node_update_verified(db, node_id, peer_name, x25519, ed_pub, join_ts); + chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id); if (db && peer_name[0]) { sqlite3_stmt* ns = NULL; sqlite3_prepare_v2(db,