Browse Source

chatgui: split chat_core into chat_conn_mgr — connection manager with rotation, max_parallel slots, background ping, per-connection DOWN callback

- chat_conn_mgr.c/h: new connection manager (1 per instance, global in g_cc)
  - cm_rotation() — single rotation method: random addr pick, INIT, retry
  - 10 parallel slots per node (CM_MAX_PARALLEL), evict worst RTT on excess
  - DEAD slot reuse before close — close only when replacement ready
  - per-connection DOWN callback via etcp_conn_add_down_cbk
  - background ping cycle: random node→random addr RTT measurement
- chat_core.c: added g_cc.conn_mgr field, init/destroy/getter
- chat_sync.c: migrate to new API, add_node() at WELCOME/JOIN/PEER_UPSERT
- Removed old ca_state/ca_ctx/ca_* single-flight code
topo_upd
Evgeny 2 months ago
parent
commit
a7d842440c
  1. 3
      tools/chatgui/CMakeLists.txt
  2. 589
      tools/chatgui/transport/chat_conn_mgr.c
  3. 55
      tools/chatgui/transport/chat_conn_mgr.h
  4. 340
      tools/chatgui/transport/chat_core.c
  5. 36
      tools/chatgui/transport/chat_core.h
  6. 47
      tools/chatgui/transport/chat_sync.c

3
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()

589
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 <sqlite3.h>
#include <string.h>
#include <stdlib.h>
#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);
}

55
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 <stdint.h>
#include <stddef.h>
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 */

340
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) {

36
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 {

47
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,

Loading…
Cancel
Save