Browse Source

refactor: move sqlite3 db from topo_group to UTUN_INSTANCE, remove conn_mgr/topo_group deps from chatgui

- topo_sqlite_db moved from TOPO_GROUPS to UTUN_INSTANCE (direct access)
- topo_groups_set_sqlite_db() removed
- chatgui: chat_sync_connect_node → chat_core_connect_auto (direct ETCP, no conn_mgr)
- chatgui: chat_core_connect_from_invite saves to SQLite → auto connect (no BGP)
- chatgui: removed _on_node_updated BGP callback from member_sync
- chatgui: removed conn_mgr_get_status from status dump
- chatgui: removed all conn_mgr.h and topo_group.h includes
- chatgui: CONN_MGR_ERR_* constants → CC_* in chat_core.h
- removed unused cc_parallel_* code, conn_state_str/conn_type_str helpers
topo_upd
Evgeny 3 months ago
parent
commit
6f4025e420
  1. 4
      src/db_sync.c
  2. 28
      src/topo_group.c
  3. 3
      src/topo_group.h
  4. 2
      src/utun_instance.h
  5. 301
      tools/chatgui/transport/chat_core.c
  6. 8
      tools/chatgui/transport/chat_core.h
  7. 40
      tools/chatgui/transport/chat_sync.c
  8. 2
      tools/chatgui/transport/chat_sync.h
  9. 37
      tools/chatgui/transport/member_sync.c
  10. 3
      tools/chatgui/transport/merkle_sync.c

4
src/db_sync.c

@ -369,8 +369,8 @@ static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t
static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t out[32])
{
if (node_id == db->inst->node_id) { memcpy(out, db->inst->my_ed25519_pubkey, 32); return 0; }
if (db->inst->topo_groups && db->inst->topo_groups->topo_sqlite_db) {
if (topo_node_sqlite_get_ed25519_pubkey(db->inst->topo_groups->topo_sqlite_db, node_id, out) == 0) return 0;
if (db->inst->topo_sqlite_db) {
if (topo_node_sqlite_get_ed25519_pubkey(db->inst->topo_sqlite_db, node_id, out) == 0) return 0;
}
struct ll_entry* e = queue_find_data_by_index(db->inst->connections, (const uint8_t*)&node_id);
if (e) {

28
src/topo_group.c

@ -254,16 +254,16 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
g->v6_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET6), "to_v6sub");
if (instance->config && instance->config->global.db_path[0]) {
int rc = sqlite3_open(instance->config->global.db_path, &g->topo_sqlite_db);
if (rc == SQLITE_OK && g->topo_sqlite_db) {
sqlite3_exec(g->topo_sqlite_db, "PRAGMA journal_mode=WAL", NULL, NULL, NULL);
sqlite3_exec(g->topo_sqlite_db, "PRAGMA foreign_keys=ON", NULL, NULL, NULL);
topo_node_sqlite_init(g->topo_sqlite_db);
int rc = sqlite3_open(instance->config->global.db_path, &instance->topo_sqlite_db);
if (rc == SQLITE_OK && instance->topo_sqlite_db) {
sqlite3_exec(instance->topo_sqlite_db, "PRAGMA journal_mode=WAL", NULL, NULL, NULL);
sqlite3_exec(instance->topo_sqlite_db, "PRAGMA foreign_keys=ON", NULL, NULL, NULL);
topo_node_sqlite_init(instance->topo_sqlite_db);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "SQLite opened: %s", instance->config->global.db_path);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "sqlite3_open(%s) failed: %s",
instance->config->global.db_path, g->topo_sqlite_db ? sqlite3_errmsg(g->topo_sqlite_db) : "null db");
if (g->topo_sqlite_db) { sqlite3_close(g->topo_sqlite_db); g->topo_sqlite_db = NULL; }
instance->config->global.db_path, instance->topo_sqlite_db ? sqlite3_errmsg(instance->topo_sqlite_db) : "null db");
if (instance->topo_sqlite_db) { sqlite3_close(instance->topo_sqlite_db); instance->topo_sqlite_db = NULL; }
}
}
@ -314,7 +314,7 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) {
memory_pool_destroy(g->v4_subnet_pool);
memory_pool_destroy(g->v6_subnet_pool);
if (g->topo_sqlite_db) { sqlite3_close(g->topo_sqlite_db); g->topo_sqlite_db = NULL; }
if (instance->topo_sqlite_db) { sqlite3_close(instance->topo_sqlite_db); instance->topo_sqlite_db = NULL; }
u_free(g); instance->topo_groups = NULL;
}
@ -340,12 +340,6 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
return group;
}
void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db) {
if (!g) return;
g->topo_sqlite_db = db;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_sqlite_db set");
}
void topo_group_set_nat_check_local(struct TOPO_GROUP* group, int allow) {
if (!group || !group->instance || !group->instance->nat_det) return;
nat_detection_set_allow_local(group->instance->nat_det, allow);
@ -587,7 +581,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) route_insert(group->instance->rt, nodeinfo1);
if (node_id != group->instance->node_id) {
sqlite3* sdb = group->instance->topo_groups->topo_sqlite_db;
sqlite3* sdb = group->instance->topo_sqlite_db;
if (sdb) {
topo_node_sqlite_node_put(sdb, nodeinfo1);
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->channel_id[0]) {
@ -646,8 +640,8 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
struct ll_entry* entry = queue_find_data_by_index(group->nodes, &node_id);
if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); }
topo_group_broadcast_withdraw(group, node_id, wd_source, sender);
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_groups->topo_sqlite_db)
topo_node_sqlite_member_del(group->instance->topo_groups->topo_sqlite_db, group->channel_id, node_id);
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_sqlite_db)
topo_node_sqlite_member_del(group->instance->topo_sqlite_db, group->channel_id, node_id);
}
return 0;
}

3
src/topo_group.h

@ -132,7 +132,6 @@ struct TOPO_GROUPS {
struct memory_pool* v6_addr_pool;
struct memory_pool* v4_subnet_pool;
struct memory_pool* v6_subnet_pool;
sqlite3* topo_sqlite_db;
topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */
};
@ -242,8 +241,6 @@ int topo_group_remove_path(struct TOPO_NODEQ* nq, struct ETCP_CONN* conn);
void topo_groups_set_node_updated_cb(struct TOPO_GROUPS* groups, topo_node_updated_fn fn);
void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db);
/** Разрешить/запретить NAT check для локальных подсетей (тестовый хелпер, делегат в nat_detection) */
void topo_group_set_nat_check_local(struct TOPO_GROUP* group, int allow);

2
src/utun_instance.h

@ -10,6 +10,7 @@ extern "C" {
#include <stdbool.h>
#include <stdio.h>
#include "../lib/memory_pool.h"
#include "../lib/sqlite3.h"
#include "secure_channel.h"
#include "etcp_api.h"
#include "config_parser.h"
@ -75,6 +76,7 @@ struct UTUN_INSTANCE {
struct ROUTE_TABLE* rt;
struct TOPO_GROUPS* topo_groups; // Groups module for topology exchange
sqlite3* topo_sqlite_db; // Shared SQLite DB (nodes/channels/peers)
struct NAT_DETECTION* nat_det; // NAT detection module
// Identification

301
tools/chatgui/transport/chat_core.c

@ -13,9 +13,6 @@
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_api.h"
#include "../../../src/topo_group.h"
#include "../../../src/topo_node.h"
#include "../../../src/conn_mgr.h"
#include "../../../src/etcp.h"
#include "../../../src/etcp_connections.h"
#include "../../../src/secure_channel.h"
@ -108,7 +105,7 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
g_cc.my_node_id = inst->node_id;
topo_node_sqlite_init(g_cc.db);
topo_groups_set_sqlite_db(inst->topo_groups, g_cc.db);
inst->topo_sqlite_db = g_cc.db;
/* записать себя в nodes + local_identity с реальным node_id и именем */
{
@ -302,9 +299,6 @@ void chat_core_update_my_name(const char* name) {
sqlite3_finalize(cs);
}
/* broadcast updated NODEINFO to connected peers */
struct TOPO_GROUP* grp = topo_groups_get_default(g_cc.inst->topo_groups);
if (grp) topo_group_update_my_nodeinfo(g_cc.inst, grp);
}
void chat_core_sync_my_addresses(void) {
@ -603,237 +597,55 @@ static void connect_result_cb(int result, uint64_t node_id, void* arg) {
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20);
}
#define CC_PARALLEL_CONNECT_TIMEOUT_MS 3000
struct cc_parallel_state {
struct UTUN_INSTANCE* inst;
int addr_count;
int pending_count;
int completed; /* 0=pending, 1=success, -1=fail-delivered */
uint64_t node_id;
uint64_t channel_id;
struct ETCP_CONN** conns;
void** timers;
};
struct cc_parallel_ctx {
struct cc_parallel_state* state;
int addr_index;
};
static void cc_parallel_cleanup(struct cc_parallel_state* st) {
if (!st) return;
u_free(st->conns);
u_free(st->timers);
u_free(st);
}
static void cc_parallel_ready_cb(struct ETCP_CONN* conn, void* arg) {
struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg;
struct cc_parallel_state* st = pctx->state;
if (st->completed) { u_free(pctx); return; }
st->completed = 1;
uint8_t data[20];
memcpy(data, &st->node_id, 8);
int r = CONN_MGR_OK;
memcpy(data + 8, &r, 4);
memcpy(data + 12, &st->channel_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect SUCCESS idx=%d peer=0x%016llx",
CC_ID, pctx->addr_index, (unsigned long long)st->node_id);
for (int i = 0; i < st->addr_count; i++) {
if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; }
if (st->conns[i] && i != pctx->addr_index) { uasync_call_soon(st->inst->ua, st->conns[i], (timeout_callback_t)etcp_connection_close); st->conns[i] = NULL; }
}
u_free(pctx);
cc_parallel_cleanup(st);
}
static void cc_parallel_timeout_cb(void* arg) {
struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg;
struct cc_parallel_state* st = pctx->state;
if (st->completed) { u_free(pctx); return; }
if (st->conns[pctx->addr_index]) { etcp_connection_close(st->conns[pctx->addr_index]); st->conns[pctx->addr_index] = NULL; }
st->timers[pctx->addr_index] = NULL;
st->pending_count--;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx",
CC_ID, pctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id);
if (st->pending_count <= 0 && !st->completed) {
st->completed = -1;
uint8_t data[20];
memcpy(data, &st->node_id, 8);
int r = CONN_MGR_ERR_TIMEOUT;
memcpy(data + 8, &r, 4);
memcpy(data + 12, &st->channel_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20);
cc_parallel_cleanup(st);
}
u_free(pctx);
}
void chat_core_connect_from_invite(struct chat_invite* inv) {
if (!g_cc.initialized || !g_cc.inst || !inv) return;
struct TOPO_GROUP* group = topo_groups_get_default(g_cc.inst->topo_groups);
if (!group || !g_cc.inst->conn_mgr) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: bgp/conn_mgr not available", CC_ID);
uint8_t err[12]; int r = -7;
memcpy(err, &inv->node_id, 8); memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12);
return;
}
uint64_t node_id = inv->node_id;
if (topo_node_find_by_id(group, node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx already in BGP, connecting",
CC_ID, (unsigned long long)node_id);
conn_mgr_connect_node(g_cc.inst->conn_mgr, node_id, 30000,
connect_result_cb, &inv->channel_id);
return;
/* save pubkey to nodes table */
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(g_cc.db,
"INSERT OR REPLACE INTO nodes(node_id, x25519_pubkey) VALUES(?,?)",
-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);
}
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_NODEQ));
if (!qe) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: failed to alloc TOPO_NODEQ", CC_ID);
uint8_t err[12]; int r = -7;
memcpy(err, &node_id, 8); memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12);
return;
}
struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)qe;
memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry));
struct TOPO_NODE* ni = u_calloc(1, sizeof(struct TOPO_NODE));
if (!ni) { queue_entry_free(qe); return; }
ni->node_id = node_id;
ni->group_id = group->group_id;
ni->ver = 0;
memcpy(ni->public_key, inv->pubkey, SC_PUBKEY_SIZE);
memset(ni->ed25519_public_key, 0, SC_PUBKEY_SIZE);
ni->node_name = u_strdup("");
struct TOPO_ADDR4* addrs_head = NULL;
const uint8_t* src = inv->addrs_data;
for (int i = 0; i < inv->addr_count; i++) {
uint8_t family = *src++;
if (family == 4) {
struct TOPO_ADDR4* a4 = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool);
if (!a4) continue;
memcpy(a4->addr, src, 4); src += 4;
a4->port = ((uint16_t)src[0] << 8) | src[1]; src += 2;
a4->type = TOPO_ADDR_REAL; a4->socket_id = 0; a4->protocol = TOPO_PROTO_UDP;
a4->next = addrs_head; addrs_head = a4;
} else {
src += 18; /* skip v6 */
/* save addresses to node_addresses */
{
sqlite3_stmt* ds = NULL;
if (sqlite3_prepare_v2(g_cc.db,
"DELETE FROM node_addresses WHERE node_id=?", -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(g_cc.db,
"INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type)"
" 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++;
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_step(is); sqlite3_reset(is);
} else {
src += 18;
}
}
sqlite3_finalize(is);
}
}
ni->v4_addrs = addrs_head;
/* alien-узел: sock_meta с неизвестной конфигурацией цели */
struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool);
if (sm) { sm->id = 0; sm->config_type = CFG_SERVER_TYPE_UNKNOWN; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = NULL; ni->v4_sock_meta = sm; }
nq->node = ni; topo_node_ref(ni);
nq->hash_node_id = node_id; nq->alien = 0;
nq->dirty = 0; nq->last_ver = 0;
nq->conn_mgr_type = CONN_TYPE_NONE; nq->conn_mgr_intermediariy_count = 0;
memset(&nq->connectivity, 0, sizeof(nq->connectivity));
nq->best_socket = NULL;
queue_data_put_with_index(group->nodes, &nq->ll);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: created NODEINFO for 0x%016llx, %d addrs, parallel connect",
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);
/* count IPv4 addresses and collect them for parallel connect */
int v4_count = 0;
for (struct TOPO_ADDR4* a = addrs_head; a; a = a->next) v4_count++;
if (v4_count == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no IPv4 addresses in invite", CC_ID);
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_NO_ADDRESSES;
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
return;
}
struct ETCP_SOCKET* best_socket = g_cc.inst->etcp_sockets;
while (best_socket && best_socket->local_addr.ss_family != AF_INET) best_socket = best_socket->next;
if (!best_socket) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no AF_INET socket", CC_ID);
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL;
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
return;
}
struct cc_parallel_state* pst = u_calloc(1, sizeof(struct cc_parallel_state));
if (!pst) {
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL;
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
return;
}
pst->inst = g_cc.inst;
pst->addr_count = v4_count;
pst->pending_count = v4_count;
pst->completed = 0;
pst->node_id = node_id;
pst->channel_id = inv->channel_id;
pst->conns = u_calloc(v4_count, sizeof(struct ETCP_CONN*));
pst->timers = u_calloc(v4_count, sizeof(void*));
if (!pst->conns || !pst->timers) { cc_parallel_cleanup(pst); return; }
int idx = 0;
for (struct TOPO_ADDR4* a = addrs_head; a && idx < v4_count; a = a->next, idx++) {
struct sockaddr_in sin;
memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa;
memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin));
struct ETCP_CONN* conn = etcp_connection_create(g_cc.inst, NULL);
if (!conn) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: etcp_connection_create failed idx=%d", CC_ID, idx);
pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue;
}
sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys);
sc_set_peer_public_key(&conn->crypto_ctx, inv->pubkey, 0);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: invite connect set peer_pubkey=%016llx my_pub=%016llx peer_node=0x%016llx",
CC_ID, *(uint64_t*)inv->pubkey, *(uint64_t*)conn->crypto_ctx.pk->public_key, inv->node_id);
struct cc_parallel_ctx* pctx = u_calloc(1, sizeof(struct cc_parallel_ctx));
if (!pctx) { etcp_connection_close(conn); pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; }
pctx->state = pst; pctx->addr_index = idx;
etcp_conn_set_ready_cbk(conn, cc_parallel_ready_cb, pctx);
if (!etcp_link_new(conn, best_socket, &sa, 0)) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: etcp_link_new failed idx=%d", CC_ID, idx);
u_free(pctx); etcp_connection_close(conn);
pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue;
}
pst->conns[idx] = conn;
etcp_conn_set_peer_node_id(conn, inv->node_id);
pst->timers[idx] = uasync_set_timeout(g_cc.inst->ua, CC_PARALLEL_CONNECT_TIMEOUT_MS * 10,
pctx, cc_parallel_timeout_cb, "cc_parallel");
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect attempt %d/%d to %d.%d.%d.%d:%d",
CC_ID, idx + 1, v4_count, sin.sin_addr.s_addr & 0xFF, (sin.sin_addr.s_addr >> 8) & 0xFF,
(sin.sin_addr.s_addr >> 16) & 0xFF, (sin.sin_addr.s_addr >> 24) & 0xFF, a->port);
}
if (pst->pending_count <= 0 && !pst->completed) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: all %d parallel connects failed to start", CC_ID, v4_count);
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_UNREACHABLE;
memcpy(err + 8, &r, 4); memcpy(err + 12, &inv->channel_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
cc_parallel_cleanup(pst);
}
chat_core_connect_auto(node_id, connect_result_cb, &inv->channel_id, NULL);
}
/* ═══ Auto-connect: direct ETCP connection from SQLite (no BGP) ═══ */
@ -882,7 +694,7 @@ static void ca_ready_cb(struct ETCP_CONN* conn, void* arg) {
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_ready_cbk(conn, ca_ready_cb, ctx);
st->result_cb(CONN_MGR_OK, st->node_id, st->result_arg);
st->result_cb(CC_OK, st->node_id, st->result_arg);
ca_cleanup(st);
}
@ -896,7 +708,7 @@ static void ca_timeout_cb(void* arg) {
CC_ID, ctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id);
if (st->pending_count <= 0 && !st->delivered) {
st->delivered = 1;
st->result_cb(CONN_MGR_ERR_TIMEOUT, st->node_id, st->result_arg);
st->result_cb(CC_ERR_TIMEOUT, st->node_id, st->result_arg);
}
}
@ -925,7 +737,7 @@ void chat_core_connect_auto(uint64_t node_id,
if (!best_socket) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no AF_INET socket for node 0x%016llx",
CC_ID, (unsigned long long)node_id);
cb(CONN_MGR_ERR_INTERNAL, node_id, arg);
cb(CC_ERR_INTERNAL, node_id, arg);
return;
}
@ -945,7 +757,7 @@ void chat_core_connect_auto(uint64_t node_id,
if (pubkey[0] == 0 && pubkey[1] == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no pubkey for node 0x%016llx",
CC_ID, (unsigned long long)node_id);
cb(CONN_MGR_ERR_NOT_FOUND, node_id, arg);
cb(CC_ERR_NOT_FOUND, node_id, arg);
return;
}
@ -978,12 +790,12 @@ void chat_core_connect_auto(uint64_t node_id,
addr_count = uniq; }
if (addr_count == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: auto_connect no IPv4 addrs yet for node 0x%016llx (BGP not synced?)", CC_ID, (unsigned long long)node_id);
cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, arg);
cb(CC_ERR_NO_ADDRESSES, node_id, arg);
return;
}
struct ca_state* pst = u_calloc(1, sizeof(struct ca_state));
if (!pst) { cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; }
if (!pst) { cb(CC_ERR_INTERNAL, node_id, arg); return; }
pst->inst = g_cc.inst;
pst->addr_count = addr_count;
pst->pending_count = addr_count;
@ -994,7 +806,7 @@ void chat_core_connect_auto(uint64_t node_id,
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) { ca_cleanup(pst); cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; }
if (!pst->conns || !pst->timers || !pst->ctxs) { ca_cleanup(pst); cb(CC_ERR_INTERNAL, node_id, arg); return; }
for (int i = 0; i < addr_count; i++) {
struct sockaddr_in sin;
@ -1039,7 +851,7 @@ void chat_core_connect_auto(uint64_t node_id,
if (pst->pending_count <= 0 && !pst->delivered) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect all %d attempts failed to start", CC_ID, addr_count);
if (out_state) *out_state = NULL;
cb(CONN_MGR_ERR_UNREACHABLE, node_id, arg);
cb(CC_ERR_UNREACHABLE, node_id, arg);
ca_cleanup(pst);
}
}
@ -1265,12 +1077,6 @@ static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, c
/* ─── сбор статуса (NTP + connections) и отправка в GUI ─── */
static const char* conn_state_str(uint8_t state) {
switch (state) { case 0: return "DISCONNECTED"; case 1: return "CONNECTING"; case 2: return "CONNECTED"; default: return "?"; }
}
static const char* conn_type_str(uint8_t t) {
switch (t) { case 0: return "NONE"; case 1: return "DIRECT"; case 2: return "REVERSE"; case 3: return "INDIRECT"; default: return "?"; }
}
static const char* nat_type_str(uint8_t t) {
switch (t) { case 0: return "UNKNOWN"; case 1: return "EIM"; case 2: return "STRICT"; case 3: return "DIRECT"; default: return "?"; }
}
@ -1335,13 +1141,6 @@ static void chat_core_collect_status(void) {
struct ETCP_LINK* tl = conn->links;
while (tl) { link_count++; tl = tl->next; }
uint8_t state = 0, ctype = 0;
conn_mgr_get_status(g_cc.inst->conn_mgr, pid, &state, &ctype);
char mgr_str[64];
if (ctype == 0) snprintf(mgr_str, sizeof(mgr_str), "%s", conn_state_str(state));
else snprintf(mgr_str, sizeof(mgr_str), "%s/%s", conn_state_str(state), conn_type_str(ctype));
const char* init_str = conn->initialized ? "" : "(!init)";
char myhex[5], peerhex[5];
@ -1352,13 +1151,13 @@ static void chat_core_collect_status(void) {
get_node_name(pid, peername, sizeof(peername));
if (peername[0])
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] %s ETCP:%s%s(%dL) MGR:%s\n",
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] %s ETCP:%s%s(%dL)\n",
myhex, peerhex, peername,
conn->links_up ? "UP" : "DOWN", init_str, link_count, mgr_str);
conn->links_up ? "UP" : "DOWN", init_str, link_count);
else
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] ETCP:%s%s(%dL) MGR:%s\n",
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] ETCP:%s%s(%dL)\n",
myhex, peerhex,
conn->links_up ? "UP" : "DOWN", init_str, link_count, mgr_str);
conn->links_up ? "UP" : "DOWN", init_str, link_count);
int link_idx = 0;
struct ETCP_LINK* link = conn->links;

8
tools/chatgui/transport/chat_core.h

@ -14,6 +14,14 @@
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);

40
tools/chatgui/transport/chat_sync.c

@ -7,8 +7,6 @@
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_api.h"
#include "../../../src/etcp.h"
#include "../../../src/conn_mgr.h"
#include "../../../src/topo_group.h"
#include "../../../src/secure_channel.h"
#include "../../../src/ntp_time.h"
#include "../../../lib/u_async.h"
@ -197,9 +195,9 @@ static void ac_fill(struct auto_connect* ac) {
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 == CONN_MGR_OK ? "OK" : "FAIL";
const char* rs = result == CC_OK ? "OK" : "FAIL";
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: result node=0x%016llx %s", AC_ID, (unsigned long long)node_id, rs);
if (result == CONN_MGR_OK) f->ca_state = NULL; /* already freed by ca_ready_cb, prevent double-free in GC/stop */
if (result == CC_OK) f->ca_state = NULL; /* already freed by ca_ready_cb, prevent double-free in GC/stop */
}
/* ── retry timer callback (every AC_RETRY_MS) ── */
@ -390,7 +388,7 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
case CS_MSG_ERROR: {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id);
if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; }
uint8_t err[20]; memcpy(err, &peer, 8); int r = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4);
uint8_t err[20]; memcpy(err, &peer, 8); int r = CC_ERR_NOT_FOUND; memcpy(err + 8, &r, 4);
memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
break;
@ -411,7 +409,7 @@ static void cs_info_req_timeout_cb(void* arg) {
CS_ID, (unsigned long long)cs->pending_invite_node_id,
(unsigned long long)cs->pending_invite_ch_id);
uint8_t err[12]; memcpy(err, &cs->pending_invite_node_id, 8);
int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12);
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
@ -425,7 +423,7 @@ static void cs_join_timeout_cb(void* arg) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &peer, 8);
int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0;
@ -522,8 +520,8 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) {
if (!conn || !g_cs) return;
uint64_t peer = conn->peer_node_id;
uint16_t rtt = conn->rtt_avg_100;
if (rtt > 0 && peer != 0 && g_cs->inst->topo_groups && g_cs->inst->topo_groups->topo_sqlite_db) {
sqlite3* db = g_cs->inst->topo_groups->topo_sqlite_db;
if (rtt > 0 && peer != 0 && g_cs->inst->topo_sqlite_db) {
sqlite3* db = g_cs->inst->topo_sqlite_db;
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db, "UPDATE node_addresses SET rtt=? WHERE node_id=?", -1, &st, NULL);
if (st) { sqlite3_bind_int(st, 1, (int)rtt); sqlite3_bind_int64(st, 2, (sqlite3_int64)peer); sqlite3_step(st); sqlite3_finalize(st); }
@ -676,12 +674,8 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
}
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
uint8_t buf[2048]; size_t buf_len;
if (chat_core_load_nodeinfo(node_id, buf, sizeof(buf), &buf_len) == 0 && buf_len > 0) {
/* nodeinfo loaded — conn_mgr will pick it up */
(void)inst;
conn_mgr_connect_node(inst->conn_mgr, node_id, 30000, NULL, NULL);
}
(void)inst;
chat_core_connect_auto(node_id, NULL, NULL, NULL);
}
void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
@ -782,7 +776,7 @@ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
(void)pl; (void)len;
char name[128]; int is_dm; uint64_t owner; uint8_t x25519[32], ed_pub[32], ch_sig[64];
if (topo_node_sqlite_channel_get(cs->inst->topo_groups->topo_sqlite_db,
if (topo_node_sqlite_channel_get(cs->inst->topo_sqlite_db,
ch_id, name, (int)sizeof(name), &is_dm, &owner, x25519, ed_pub, ch_sig) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: CHANNEL_INFO_REQ unknown ch=%s from=%016llx", CS_ID, ch_id, (unsigned long long)peer);
uint8_t err[1] = { CS_MSG_ERROR }; cs_send(cs, ch_id, peer, err, 1);
@ -795,7 +789,7 @@ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
/* load or create join_sig */
uint8_t my_join_sig[64] = {0}; uint64_t my_join_ts = 0;
{
sqlite3* vdb = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* vdb = cs->inst->topo_sqlite_db;
if (topo_node_sqlite_member_get_join(vdb, ch_id, myid, my_join_sig, &my_join_ts) != 0) {
uint8_t join_msg[256]; size_t mlen = 0;
mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", ch_id) + 1;
@ -886,7 +880,7 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
}
/* save channel to local DB */
topo_node_sqlite_channel_put(cs->inst->topo_groups->topo_sqlite_db,
topo_node_sqlite_channel_put(cs->inst->topo_sqlite_db,
ch_id, name, (int)is_dm, owner, x25519, NULL, ed_pub, NULL, ch_sig);
chat_core_ensure_channel_ready(ch_id);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready for sync ch=%s, db_sync will pick up via periodic check or active conn",
@ -914,7 +908,7 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
CS_ID, (unsigned long long)peer, (unsigned long long)inviter_update_ts);
}
/* save inviter node_info to local DB */
sqlite3* vdb = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* vdb = cs->inst->topo_sqlite_db;
if (vdb && inv_name[0]) {
sqlite3_stmt* ns = NULL;
sqlite3_prepare_v2(vdb, "INSERT OR REPLACE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey) VALUES(?,?,?,?)", -1, &ns, NULL);
@ -1047,7 +1041,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
return;
}
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* db = cs->inst->topo_sqlite_db;
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_CONNECTIVITY, "%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);
@ -1143,7 +1137,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; }
(void)peer;
if (len < 2) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: WELCOME too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; }
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* db = cs->inst->topo_sqlite_db;
const uint8_t* p = pl;
uint16_t pc; memcpy(&pc, p, 2); p += 2;
@ -1251,7 +1245,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
return;
}
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* db = cs->inst->topo_sqlite_db;
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_CONNECTIVITY, "%s: member_put(peer_upsert) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
@ -1305,7 +1299,7 @@ static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: PEER_REMOVE too short len=%zu", CS_ID, len); return; }
uint64_t node_id; memcpy(&node_id, pl, 8);
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db;
sqlite3* db = cs->inst->topo_sqlite_db;
topo_node_sqlite_member_del(db, ch_id, node_id);
cs_refresh_channels(cs);

2
tools/chatgui/transport/chat_sync.h

@ -46,7 +46,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst,
void (*gui_result_cb)(void*, int, const uint8_t*, int));
void chat_sync_destroy(struct UTUN_INSTANCE* inst);
/* Прослойка conn_mgr: загрузить nodeinfo из БД и подключиться */
/* Прямое ETCP-подключение к узлу (pubkey+адреса из БД) */
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Подключиться к пиру по данным invite-ссылки (вызывается из GUI-потока) */

37
tools/chatgui/transport/member_sync.c

@ -3,7 +3,6 @@
#include "chat_core.h"
#include "../../../src/utun_instance.h"
#include "../../../src/topo_group.h"
#include "../../../lib/debug_config.h"
#include "../../../lib/mem.h"
@ -19,7 +18,7 @@ struct addr_item { uint8_t family; uint8_t addr[16]; uint16_t port; };
/* ── DB access ── */
static sqlite3* _db(struct UTUN_INSTANCE* inst) {
return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL;
return inst ? inst->topo_sqlite_db : NULL;
}
static void _peers_table(const char* ch_id, char* buf, size_t sz) {
@ -281,50 +280,18 @@ static const struct merkle_sync_data_ops g_member_ops = {
.apply_items = _member_apply_items,
};
/* ── node_updated callback ── */
static void _on_node_updated(struct UTUN_INSTANCE* inst, uint64_t node_id,
const uint8_t* x25519, const uint8_t* ed25519) {
(void)x25519; (void)ed25519;
if (!inst) return;
if (node_id == inst->node_id) chat_core_sync_my_addresses();
sqlite3* db = _db(inst); if (!db) return;
int rc = 0;
sqlite3_stmt* cs = NULL;
if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return;
while (sqlite3_step(cs) == SQLITE_ROW) {
const char* ch = (const char*)sqlite3_column_text(cs, 0);
if (!ch) continue;
char peers_tbl[128]; _peers_table(ch, peers_tbl, sizeof(peers_tbl));
char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl);
sqlite3_stmt* ps = NULL;
if (sqlite3_prepare_v2(db, buf, -1, &ps, NULL) == SQLITE_OK) {
sqlite3_bind_int64(ps, 1, (sqlite3_int64)node_id);
if (sqlite3_step(ps) == SQLITE_ROW) {
merkle_sync_recompute_path(inst, ch, node_id);
rc++;
}
sqlite3_finalize(ps);
}
}
sqlite3_finalize(cs);
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: _on_node_updated node=%016llx channels_recomputed=%d", MS_ID, (unsigned long long)node_id, rc);
}
/* ── Public API ── */
int member_sync_init(struct UTUN_INSTANCE* inst) {
if (!inst) return -1;
int rc = merkle_sync_init(inst, 0x31, &g_member_ops, inst);
if (rc != 0) return rc;
topo_groups_set_node_updated_cb(inst->topo_groups, _on_node_updated);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized, merkle_rc=%d cb_registered=%d", MS_ID, rc, inst->topo_groups && inst->topo_groups->node_updated_cb ? 1 : 0);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized, merkle_rc=%d", MS_ID, rc);
return 0;
}
void member_sync_destroy(struct UTUN_INSTANCE* inst) {
if (!inst) return;
topo_groups_set_node_updated_cb(inst->topo_groups, NULL);
merkle_sync_destroy(inst);
}

3
tools/chatgui/transport/merkle_sync.c

@ -3,7 +3,6 @@
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_api.h"
#include "../../../src/etcp.h"
#include "../../../src/topo_group.h"
#include "../../../lib/debug_config.h"
#include "../../../lib/mem.h"
#include "../../../lib/u_async.h"
@ -44,7 +43,7 @@ struct merkle_sync {
/* ── DB access ── */
static sqlite3* _db(struct UTUN_INSTANCE* inst) {
return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL;
return inst ? inst->topo_sqlite_db : NULL;
}
/* ── Prefix arithmetic ── */

Loading…
Cancel
Save