Browse Source

topo_group_connect: авто-подключение CHAT-групп + conn_mgr_connect_from_invite + замена chat_conn_mgr в chatgui

Добавлено:
- conn_mgr_connect_from_invite() — подключение к новому узлу по invite-данным
  (INIT → INVITE_INFO_REQ/RESP → проверка членства в группе → сохранение в БД)
- topo_group_connect.c — авто-подключение при старте CHAT-группы:
  Phase 1: одновременный connect всех connected=1 пиров, таймаут, отмена незавершённых
  Phase 2: последовательный перебор публичных адресов до первого успеха
- Счётчик active_conn_count с проверкой nodeinfo-путей (UP/DOWN)
- restart() при phase=DONE и count==0

БД:
- connected INTEGER NOT NULL DEFAULT 0 в peers_<ch> + ALTER TABLE для существующих
- topo_node_sqlite_set_connected / get_connected_peers / get_public_peers

Chatgui:
- Включён topo_group (убран utun_instance_set_topo_group_enabled(0))
- Создание TOPO_GROUP_TYPE_CHAT при channel_ready и при старте
- Invite connect через conn_mgr_connect_from_invite
- Удалён chat_conn_mgr (из сборки) + весь старый auto-connect код
topo_upd
Evgeny 2 months ago
parent
commit
6edd40d141
  1. 2
      src/Makefile.am
  2. 289
      src/routing_layer/conn_mgr.c
  3. 55
      src/routing_layer/conn_mgr.h
  4. 6
      src/routing_layer/topo_group.c
  5. 8
      src/routing_layer/topo_group.h
  6. 260
      src/routing_layer/topo_group_connect.c
  7. 30
      src/routing_layer/topo_group_connect.h
  8. 63
      src/routing_layer/topo_node_sqlite.c
  9. 4
      src/routing_layer/topo_node_sqlite.h
  10. 3
      tools/chatgui/CMakeLists.txt
  11. 4
      tools/chatgui/libutun/CMakeLists.txt
  12. 16
      tools/chatgui/transport/chat_core.c
  13. 2
      tools/chatgui/transport/chat_core.h
  14. 333
      tools/chatgui/transport/chat_sync.c
  15. 5
      tools/chatgui/transport/chat_sync.h
  16. 25
      tools/chatgui/transport/utun_node.cpp

2
src/Makefile.am

@ -17,6 +17,7 @@ utun_CORE_SOURCES = \
routing_layer/route_connectivity.c \ routing_layer/route_connectivity.c \
routing_layer/conn_mgr.c \ routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \ routing_layer/topo_recovery.c \
routing_layer/topo_group_connect.c \
db_sync.c \ db_sync.c \
routing_layer/routing.c \ routing_layer/routing.c \
tun_if.c \ tun_if.c \
@ -73,6 +74,7 @@ libutun_a_SOURCES = \
routing_layer/route_connectivity.c \ routing_layer/route_connectivity.c \
routing_layer/conn_mgr.c \ routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \ routing_layer/topo_recovery.c \
routing_layer/topo_group_connect.c \
db_sync.c \ db_sync.c \
routing_layer/routing.c \ routing_layer/routing.c \
tun_if.c \ tun_if.c \

289
src/routing_layer/conn_mgr.c

@ -54,6 +54,14 @@ static void cm_reverse_timeout_cb(void* arg);
static void cm_exchange_timeout_cb(void* arg); static void cm_exchange_timeout_cb(void* arg);
static void cm_candidate_ping_timer_cb(void* arg); static void cm_candidate_ping_timer_cb(void* arg);
static void cm_invite_init_cb(struct ETCP_CONN* conn, void* arg);
static void cm_invite_overall_timeout(void* arg);
static void cm_invite_fail(struct cm_invite_pending* inv, int result);
static void cm_invite_cleanup(struct cm_invite_pending* inv);
static void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn,
const uint8_t* data, size_t len);
static void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len);
struct cm_reverse_pending { struct cm_reverse_pending {
struct cm_reverse_pending* next; struct cm_reverse_pending* next;
uint32_t request_id; uint32_t request_id;
@ -72,6 +80,27 @@ struct cm_exchange_pending {
uint8_t probes_done; uint8_t probes_done;
}; };
#define CM_INVITE_DEFAULT_TIMEOUT_MS 2000
enum CM_INVITE_STATE {
CM_INVITE_CONNECTING = 0,
CM_INVITE_WAIT_INFO = 1,
};
struct cm_invite_pending {
struct cm_invite_pending* next;
struct CONN_MGR* mgr;
uint64_t node_id;
uint64_t group_id;
uint32_t timeout_ms;
uint8_t state;
struct ETCP_CONN* conn;
void* overall_timer;
conn_mgr_connect_callback_t cb;
void* cb_arg;
struct TOPO_NODEQ* temp_nq;
};
struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) { struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) {
if (!group || !group->instance) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NULL group"); return NULL; } if (!group || !group->instance) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NULL group"); return NULL; }
struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR)); struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR));
@ -82,7 +111,9 @@ struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) {
mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY)); mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY));
if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "entries alloc failed"); u_free(mgr); return NULL; } if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "entries alloc failed"); u_free(mgr); return NULL; }
mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS; mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS;
mgr->invite_list = NULL;
etcp_router_bind(mgr->instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler); etcp_router_bind(mgr->instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler);
etcp_bind(mgr->instance, ETCP_RT_ID_CONN_MGR, conn_mgr_router_recv_handler);
mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping"); mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping");
mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping");
mgr->bg_ping_cycle_start_tb = get_time_tb(); mgr->bg_ping_cycle_start_tb = get_time_tb();
@ -96,7 +127,9 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) {
if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; } if (mgr->bg_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->bg_ping_timer); mgr->bg_ping_timer = NULL; }
if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; } if (mgr->candidate_ping_timer) { uasync_cancel_timeout(mgr->instance->ua, mgr->candidate_ping_timer); mgr->candidate_ping_timer = NULL; }
etcp_router_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR); etcp_router_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR);
etcp_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR);
for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]); for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_destroy(&mgr->entries[i]);
while (mgr->invite_list) cm_invite_fail(mgr->invite_list, CONN_MGR_ERR_INTERNAL);
struct cm_exchange_pending* ep = mgr->exchange_pending; struct cm_exchange_pending* ep = mgr->exchange_pending;
while (ep) { while (ep) {
struct cm_exchange_pending* next = ep->next; struct cm_exchange_pending* next = ep->next;
@ -914,6 +947,13 @@ static void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry
cm_handle_disconnect(mgr, disc->node_id); cm_handle_disconnect(mgr, disc->node_id);
} }
break; break;
case CONN_MGR_SUBCMD_INVITE_INFO_REQ:
if (mgr && len >= sizeof(struct CONN_MGR_INVITE_INFO_REQ))
cm_handle_invite_info_req(mgr, conn, dgram, len);
break;
case CONN_MGR_SUBCMD_INVITE_INFO_RESP:
if (mgr) cm_handle_invite_info_resp(mgr, dgram, len);
break;
} }
} }
@ -1097,3 +1137,252 @@ static void cm_candidate_ping_timer_cb(void* arg) {
mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB,
mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping"); mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping");
} }
/* ═══════════════════════════════════════════════════════════════════════
* Invite-подключение
* ══════════════════════════════════════════════════════════════════════ */
static void sanitize_ch_id(const char* ch_id, char* out, size_t out_sz) {
size_t j = 0;
for (size_t i = 0; ch_id[i] && j < out_sz - 1; i++)
if ((ch_id[i] >= '0' && ch_id[i] <= '9') || (ch_id[i] >= 'a' && ch_id[i] <= 'z') || (ch_id[i] >= 'A' && ch_id[i] <= 'Z') || ch_id[i] == '_')
out[j++] = ch_id[i];
out[j] = '\0';
}
static int cm_check_group_member(struct CONN_MGR* mgr, uint64_t group_id, uint64_t peer_node_id) {
struct TOPO_GROUP* group = topo_groups_find(mgr->instance->topo_groups, group_id);
if (!group) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite_info group 0x%016llx not found", (unsigned long long)group_id); return 0; }
if (group->group_type == TOPO_GROUP_TYPE_CHAT) {
sqlite3* db = mgr->instance->topo_sqlite_db;
if (!db || !group->channel_id[0]) return 0;
char san[64]; sanitize_ch_id(group->channel_id, san, sizeof(san));
char tbl[80]; snprintf(tbl, sizeof(tbl), "peers_%s", san);
char sql[128]; snprintf(sql, sizeof(sql), "SELECT 1 FROM \"%s\" WHERE node_id=?", tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id);
int is_member = (sqlite3_step(st) == SQLITE_ROW);
sqlite3_finalize(st);
return is_member;
}
return 0;
}
return topo_node_find_by_id(group, peer_node_id) != NULL;
}
static void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn,
const uint8_t* data, size_t len) {
if (len < sizeof(struct CONN_MGR_INVITE_INFO_REQ)) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ too short (%zu)", len);
return;
}
struct CONN_MGR_INVITE_INFO_REQ* req = (struct CONN_MGR_INVITE_INFO_REQ*)data;
uint64_t peer_node_id = conn->peer_node_id;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ from 0x%016llx group=0x%016llx",
(unsigned long long)peer_node_id, (unsigned long long)req->group_id);
uint8_t is_member = 0;
if (req->group_id != 0) is_member = (uint8_t)cm_check_group_member(mgr, req->group_id, peer_node_id);
struct TOPO_GROUP* group = topo_groups_find(mgr->instance->topo_groups, req->group_id);
struct TOPO_GROUP* src_group = group ? group : mgr->group;
const uint8_t* ed_pubkey = src_group->ed25519_public_key;
const char* node_name = (src_group->local_node && src_group->local_node->node)
? src_group->local_node->node->node_name : NULL;
size_t name_len = node_name ? strlen(node_name) : 0;
if (name_len > 255) name_len = 255;
size_t resp_size = sizeof(struct CONN_MGR_INVITE_INFO_RESP) + name_len;
struct CONN_MGR_INVITE_INFO_RESP* resp = u_calloc(1, resp_size);
if (!resp) return;
resp->cmd = ETCP_RT_ID_CONN_MGR;
resp->subcmd = CONN_MGR_SUBCMD_INVITE_INFO_RESP;
resp->group_id = req->group_id;
resp->is_member = is_member;
memcpy(resp->ed25519_pubkey, ed_pubkey, SC_PUBKEY_SIZE);
resp->node_name_len = (uint8_t)name_len;
if (name_len) memcpy(resp->node_name, node_name, name_len);
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = (uint8_t*)resp; qe->len = (uint16_t)resp_size; etcp_send(conn, qe); }
else u_free(resp);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_RESP to 0x%016llx member=%d name_len=%zu",
(unsigned long long)peer_node_id, is_member, name_len);
}
static void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
if (len < sizeof(struct CONN_MGR_INVITE_INFO_RESP)) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP too short (%zu)", len);
return;
}
struct CONN_MGR_INVITE_INFO_RESP* resp = (struct CONN_MGR_INVITE_INFO_RESP*)data;
size_t expected = sizeof(struct CONN_MGR_INVITE_INFO_RESP) + (size_t)resp->node_name_len;
if (len < expected) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP truncated (%zu < %zu)", len, expected);
return;
}
struct cm_invite_pending* prev = NULL;
struct cm_invite_pending* inv = mgr->invite_list;
while (inv) {
if (inv->node_id == resp->group_id ? 0 : 0) { inv = inv->next; continue; }
/* find by state == CM_INVITE_WAIT_INFO and node_id match via conn */
if (inv->state == CM_INVITE_WAIT_INFO && inv->conn) {
break;
}
prev = inv; inv = inv->next;
}
/* need to find the right invite by peer_node_id from the connection that delivered this */
struct cm_invite_pending* found = NULL;
for (struct cm_invite_pending* p = mgr->invite_list; p; p = p->next)
if (p->state == CM_INVITE_WAIT_INFO) { found = p; break; }
if (!found) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP — no pending invite in WAIT_INFO state"); return; }
inv = found;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP node=0x%016llx group=0x%016llx member=%d name_len=%u",
(unsigned long long)inv->node_id, (unsigned long long)resp->group_id, resp->is_member, resp->node_name_len);
if (inv->group_id != 0 && !resp->is_member) { cm_invite_fail(inv, CONN_MGR_ERR_NOT_MEMBER); return; }
struct TOPO_NODEQ* nq = inv->temp_nq;
if (nq && nq->node) {
memcpy(nq->node->ed25519_public_key, resp->ed25519_pubkey, SC_PUBKEY_SIZE);
if (nq->node->node_name) { u_free(nq->node->node_name); nq->node->node_name = NULL; }
if (resp->node_name_len) { nq->node->node_name = u_strdup((const char*)resp->node_name); }
sqlite3* db = inv->mgr->instance->topo_sqlite_db;
if (db) topo_node_sqlite_node_put(db, nq, (time_t)(get_time_tb() / 10000));
}
if (inv->overall_timer) { uasync_cancel_timeout(inv->mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; }
if (inv->cb) inv->cb(CONN_MGR_OK, inv->node_id, inv->cb_arg);
cm_invite_cleanup(inv);
}
static void cm_invite_init_cb(struct ETCP_CONN* conn, void* arg) {
struct cm_invite_pending* inv = (struct cm_invite_pending*)arg;
if (!conn || conn->peer_node_id != inv->node_id) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite init cb peer mismatch conn=0x%016llx inv=0x%016llx",
(unsigned long long)(conn ? conn->peer_node_id : 0), (unsigned long long)inv->node_id);
if (conn) cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite INIT OK node=0x%016llx", (unsigned long long)inv->node_id);
inv->state = CM_INVITE_WAIT_INFO;
struct CONN_MGR_INVITE_INFO_REQ req;
req.cmd = ETCP_RT_ID_CONN_MGR; req.subcmd = CONN_MGR_SUBCMD_INVITE_INFO_REQ;
req.group_id = inv->group_id;
struct ll_entry* qe = queue_entry_new(0);
if (!qe) { cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL); return; }
qe->dgram = u_malloc(sizeof(req));
if (!qe->dgram) { queue_entry_free(qe); cm_invite_fail(inv, CONN_MGR_ERR_INTERNAL); return; }
memcpy(qe->dgram, &req, sizeof(req)); qe->len = sizeof(req);
etcp_send(conn, qe);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_REQ to 0x%016llx group=0x%016llx",
(unsigned long long)inv->node_id, (unsigned long long)inv->group_id);
}
static void cm_invite_overall_timeout(void* arg) {
struct cm_invite_pending* inv = (struct cm_invite_pending*)arg;
inv->overall_timer = NULL;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite overall timeout node=0x%016llx", (unsigned long long)inv->node_id);
cm_invite_fail(inv, CONN_MGR_ERR_TIMEOUT);
}
static void cm_invite_fail(struct cm_invite_pending* inv, int result) {
if (!inv) return;
struct CONN_MGR* mgr = inv->mgr;
if (inv->overall_timer) { uasync_cancel_timeout(mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; }
if (inv->conn) { etcp_connection_close(inv->conn); inv->conn = NULL; }
if (inv->temp_nq) {
struct TOPO_GROUP* group = mgr->group;
queue_remove_data(group->nodes, &inv->temp_nq->ll);
topo_node_free_lists(group, inv->temp_nq);
queue_entry_free(&inv->temp_nq->ll); inv->temp_nq = NULL;
}
struct cm_invite_pending** pp = &mgr->invite_list;
while (*pp) { if (*pp == inv) { *pp = inv->next; break; } pp = &(*pp)->next; }
if (inv->cb) inv->cb(result, inv->node_id, inv->cb_arg);
cm_invite_cleanup(inv);
}
static void cm_invite_cleanup(struct cm_invite_pending* inv) {
if (!inv) return;
u_free(inv);
}
int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
uint64_t group_id, uint32_t timeout_ms,
conn_mgr_connect_callback_t cb, void* cb_arg) {
if (!mgr || !ni || !ni->v4_addrs) { if (cb) cb(CONN_MGR_ERR_INTERNAL, ni ? ni->node_id : 0, cb_arg); return CONN_MGR_ERR_INTERNAL; }
uint64_t node_id = ni->node_id;
if (node_id == 0 || node_id == mgr->instance->node_id) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
struct TOPO_GROUP* group = mgr->group;
struct TOPO_GROUPS* groups = mgr->instance->topo_groups;
if (!groups) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
if (topo_node_find_by_id(group, node_id)) { if (cb) cb(CONN_MGR_ERR_ALREADY_CONNECTED, node_id, cb_arg); return CONN_MGR_ERR_ALREADY_CONNECTED; }
if (!timeout_ms) timeout_ms = CM_INVITE_DEFAULT_TIMEOUT_MS;
struct TOPO_NODE* group_ni = u_calloc(1, sizeof(struct TOPO_NODE));
if (!group_ni) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
group_ni->node_id = node_id; group_ni->ver = 1;
memcpy(group_ni->public_key, ni->public_key, SC_PUBKEY_SIZE);
for (const struct TOPO_ADDR4* src = ni->v4_addrs; src; src = src->next) {
struct TOPO_ADDR4* a4 = memory_pool_alloc(groups->v4_addr_pool);
if (!a4) continue;
memset(a4, 0, sizeof(*a4)); memcpy(a4->addr, src->addr, 4); a4->port = src->port;
a4->type = TOPO_ADDR_NAT; a4->protocol = TOPO_PROTO_UDP; a4->socket_id = 0;
a4->next = group_ni->v4_addrs; group_ni->v4_addrs = a4;
}
if (!group_ni->v4_addrs) { u_free(group_ni); if (cb) cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, cb_arg); return CONN_MGR_ERR_NO_ADDRESSES; }
group_ni = topo_node_registry_acquire(groups, group_ni); /* group_ni now owned by registry (ref=1) */
if (!group_ni) { if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_NODEQ));
if (!qe) { topo_node_registry_release(groups, group_ni); if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)qe;
memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry));
nq->node = group_ni; nq->hash_node_id = node_id;
struct cm_invite_pending* inv = u_calloc(1, sizeof(struct cm_invite_pending));
if (!inv) { queue_entry_free(qe); topo_node_registry_release(groups, group_ni); if (cb) cb(CONN_MGR_ERR_INTERNAL, node_id, cb_arg); return CONN_MGR_ERR_INTERNAL; }
inv->mgr = mgr; inv->node_id = node_id; inv->group_id = group_id;
inv->timeout_ms = timeout_ms; inv->state = CM_INVITE_CONNECTING;
inv->temp_nq = nq; inv->cb = cb; inv->cb_arg = cb_arg;
inv->next = mgr->invite_list; mgr->invite_list = inv;
queue_data_put_with_index(group->nodes, &nq->ll);
int connected = 0;
for (const struct TOPO_ADDR4* a = group_ni->v4_addrs; a; a = a->next) {
if (a->protocol != TOPO_PROTO_UDP) continue;
struct ETCP_SOCKET* s = mgr->instance->etcp_sockets;
while (s) {
if (s->local_addr.ss_family == AF_INET) {
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
struct sockaddr_storage sa; memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin));
struct ETCP_CONN* conn = etcp_connection_create(mgr->instance, NULL);
if (conn) {
etcp_conn_add_init_cbk(conn, cm_invite_init_cb, inv);
sc_init_ctx(&conn->crypto_ctx, &mgr->instance->my_keys);
sc_set_peer_public_key(&conn->crypto_ctx, group_ni->public_key, 0);
if (etcp_link_new(conn, s, &sa, 0)) { inv->conn = conn; connected = 1; break; }
etcp_connection_close(conn);
}
}
s = s->next;
}
if (connected) break;
}
if (!connected) { cm_invite_fail(inv, CONN_MGR_ERR_UNREACHABLE); return CONN_MGR_ERR_UNREACHABLE; }
inv->overall_timer = uasync_set_timeout(mgr->instance->ua, (int)inv->timeout_ms * 10, inv, cm_invite_overall_timeout, "conn_mgr_invite");
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite started node=0x%016llx group=0x%016llx timeout=%ums",
(unsigned long long)node_id, (unsigned long long)group_id, timeout_ms);
return CONN_MGR_OK;
}

55
src/routing_layer/conn_mgr.h

@ -83,6 +83,7 @@ struct TOPO_GROUP;
#define CONN_MGR_ERR_REFUSED -5 ///< Соединение отклонено (зарезервировано) #define CONN_MGR_ERR_REFUSED -5 ///< Соединение отклонено (зарезервировано)
#define CONN_MGR_ERR_ALREADY_CONNECTED -6 ///< Уже подключены #define CONN_MGR_ERR_ALREADY_CONNECTED -6 ///< Уже подключены
#define CONN_MGR_ERR_INTERNAL -7 ///< Внутренняя ошибка #define CONN_MGR_ERR_INTERNAL -7 ///< Внутренняя ошибка
#define CONN_MGR_ERR_NOT_MEMBER -8 ///< Узел не является членом группы
/* ==================== Состояния ==================== */ /* ==================== Состояния ==================== */
@ -119,6 +120,8 @@ enum CM_TRY_STATE {
#define CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP 0x04 ///< Ответ с нашими кандидатами + замеры RTT до кандидатов инициатора #define CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP 0x04 ///< Ответ с нашими кандидатами + замеры RTT до кандидатов инициатора
#define CONN_MGR_SUBCMD_INTERM_SELECTED 0x05 ///< Выбранные посредники (финальный шаг INDIRECT-фазы) #define CONN_MGR_SUBCMD_INTERM_SELECTED 0x05 ///< Выбранные посредники (финальный шаг INDIRECT-фазы)
#define CONN_MGR_SUBCMD_DISCONNECT 0x06 ///< Уведомление о разрыве соединения #define CONN_MGR_SUBCMD_DISCONNECT 0x06 ///< Уведомление о разрыве соединения
#define CONN_MGR_SUBCMD_INVITE_INFO_REQ 0x0A ///< Запрос проверки членства в группе при invite-подключении
#define CONN_MGR_SUBCMD_INVITE_INFO_RESP 0x0B ///< Ответ с членством, ed25519_pubkey, именем узла
/** @} */ /** @} */
/** /**
@ -220,6 +223,38 @@ struct CONN_MGR_DISCONNECT {
#pragma pack(pop) #pragma pack(pop)
/**
* @struct CONN_MGR_INVITE_INFO_REQ
* @brief Запрос проверки членства в группе (invite-подключение).
*
* Отправляется напрямую по ETCP после INIT-рукопожатия.
* Запрашивает: ed25519_pubkey, имя узла, членство в группе.
*/
#pragma pack(push, 1)
struct CONN_MGR_INVITE_INFO_REQ {
uint8_t cmd; ///< ETCP_RT_ID_CONN_MGR
uint8_t subcmd; ///< CONN_MGR_SUBCMD_INVITE_INFO_REQ
uint64_t group_id; ///< целевая группа для проверки членства
};
#pragma pack(pop)
/**
* @struct CONN_MGR_INVITE_INFO_RESP
* @brief Ответ на INVITE_INFO_REQ.
*/
#pragma pack(push, 1)
struct CONN_MGR_INVITE_INFO_RESP {
uint8_t cmd; ///< ETCP_RT_ID_CONN_MGR
uint8_t subcmd; ///< CONN_MGR_SUBCMD_INVITE_INFO_RESP
uint64_t group_id; ///< эхо group_id из запроса
uint8_t is_member; ///< 1 = член группы, 0 = нет
uint8_t reserved[3];
uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; ///< Ed25519 публичный ключ
uint8_t node_name_len; ///< длина имени (0..255)
uint8_t node_name[0]; ///< UTF-8 имя (variable-length)
};
#pragma pack(pop)
/* ==================== Размеры пакетов протокола ==================== */ /* ==================== Размеры пакетов протокола ==================== */
#define CONN_MGR_DIRECT_REQ_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_REQ) ///< Размер заголовка DIRECT_REQ (без addr_count адресов) #define CONN_MGR_DIRECT_REQ_HDR_SIZE sizeof(struct CONN_MGR_DIRECT_REQ) ///< Размер заголовка DIRECT_REQ (без addr_count адресов)
@ -228,6 +263,7 @@ struct CONN_MGR_DISCONNECT {
#define CONN_MGR_INTERM_EXCHANGE_RESP_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_RESP) ///< Размер INTERM_EXCHANGE_RESP (макс.) #define CONN_MGR_INTERM_EXCHANGE_RESP_SIZE sizeof(struct CONN_MGR_INTERM_EXCHANGE_RESP) ///< Размер INTERM_EXCHANGE_RESP (макс.)
#define CONN_MGR_INTERM_SELECTED_SIZE sizeof(struct CONN_MGR_INTERM_SELECTED) ///< Размер INTERM_SELECTED #define CONN_MGR_INTERM_SELECTED_SIZE sizeof(struct CONN_MGR_INTERM_SELECTED) ///< Размер INTERM_SELECTED
#define CONN_MGR_DISCONNECT_SIZE sizeof(struct CONN_MGR_DISCONNECT) ///< Размер DISCONNECT #define CONN_MGR_DISCONNECT_SIZE sizeof(struct CONN_MGR_DISCONNECT) ///< Размер DISCONNECT
#define CONN_MGR_INVITE_INFO_REQ_SIZE sizeof(struct CONN_MGR_INVITE_INFO_REQ) ///< Размер INVITE_INFO_REQ
/* ==================== Callback-типы и внутренние структуры ==================== */ /* ==================== Callback-типы и внутренние структуры ==================== */
@ -319,6 +355,7 @@ struct CONN_MGR {
struct cm_reverse_pending* reverse_pending; ///< Связный список ожидающих REVERSE-подключений struct cm_reverse_pending* reverse_pending; ///< Связный список ожидающих REVERSE-подключений
struct cm_exchange_pending* exchange_pending; ///< Связный список ожидающих INDIRECT-обменов struct cm_exchange_pending* exchange_pending; ///< Связный список ожидающих INDIRECT-обменов
struct cm_invite_pending* invite_list; ///< Связный список ожидающих invite-подключений
}; };
/* ==================== Жизненный цикл ==================== */ /* ==================== Жизненный цикл ==================== */
@ -411,6 +448,24 @@ int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t t
*/ */
void conn_mgr_set_direct_timeout_ms(struct CONN_MGR* mgr, uint32_t timeout_ms); void conn_mgr_set_direct_timeout_ms(struct CONN_MGR* mgr, uint32_t timeout_ms);
/**
* @brief Подключается к новому узлу по invite-данным и проверяет членство в группе.
* @param mgr менеджер группы
* @param ni данные узла (caller-owned, только для чтения): node_id, public_key, v4_addrs
* @param group_id группа для проверки членства (0 = без проверки)
* @param timeout_ms общий таймаут (0 = 2000ms по умолчанию)
* @param cb callback результата
* @param cb_arg аргумент для callback
* @return CONN_MGR_OK если процесс запущен, CONN_MGR_ERR_* при ошибке
*
* Поток: DIRECT-подключение → INIT → INVITE_INFO_REQ → INVITE_INFO_RESP → успех/ошибка.
* При успехе узел добавляется в group->nodes и сохраняется в SQLite.
* Callback гарантированно вызывается (успех, таймаут, ошибка членства).
*/
int conn_mgr_connect_from_invite(struct CONN_MGR* mgr, struct TOPO_NODE* ni,
uint64_t group_id, uint32_t timeout_ms,
conn_mgr_connect_callback_t cb, void* cb_arg);
/** /**
* @brief Получает текущий статус подключения к ноде. * @brief Получает текущий статус подключения к ноде.
* @param mgr менеджер * @param mgr менеджер

6
src/routing_layer/topo_group.c

@ -23,6 +23,7 @@
#include "route_connectivity.h" #include "route_connectivity.h"
#include "control_server.h" #include "control_server.h"
#include "conn_mgr.h" #include "conn_mgr.h"
#include "topo_group_connect.h"
#include "topo_recovery.h" #include "topo_recovery.h"
@ -238,6 +239,7 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group_id=%016llx", (unsigned long long)group->group_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "group_id=%016llx", (unsigned long long)group->group_id);
topo_recovery_cancel_all(group); topo_recovery_cancel_all(group);
topo_group_connect_destroy(group);
struct ll_entry* e; struct ll_entry* e;
while ((e = queue_data_get(group->senders_list)) != NULL) queue_entry_free(e); while ((e = queue_data_get(group->senders_list)) != NULL) queue_entry_free(e);
@ -381,6 +383,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
} }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Group created: group_id=%016llx type=%u ch_id=%s", (unsigned long long)group_id, group_type, group->channel_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "Group created: group_id=%016llx type=%u ch_id=%s", (unsigned long long)group_id, group_type, group->channel_id);
if (group_type == TOPO_GROUP_TYPE_CHAT) topo_group_connect_init(group);
return group; return group;
} }
@ -406,6 +409,7 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id); DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id);
topo_group_send_table_request(group, conn); topo_group_send_table_request(group, conn);
topo_group_connect_on_up(group, conn);
} }
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
@ -451,6 +455,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (cascaded > 0) topo_recovery_start(group); if (cascaded > 0) topo_recovery_start(group);
topo_group_connect_on_down(group, conn);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d", conn->log_name, nodes_removed); DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d", conn->log_name, nodes_removed);
} }

8
src/routing_layer/topo_group.h

@ -2,6 +2,12 @@
* @file topo_group.h * @file topo_group.h
* @brief BGP-подобный обмен топологией узлов между пирами через ETCP. * @brief BGP-подобный обмен топологией узлов между пирами через ETCP.
* *
* Идея: модуль автоматически выстраивает карту маршрутизации между узлами,
* используя только пассивное наблюдение (подписка на события)
* Можно иметь много групп узлов. Группа - это список узлов.
* Для каждой группы выстраивается своя независимая таблица маршрутизации.
* Один узел может входить в любое число групп.
* Механика: * Механика:
* - Узлы обмениваются информацией друг о друге (pubkey, адреса, подсети) * - Узлы обмениваются информацией друг о друге (pubkey, адреса, подсети)
* через NODEINFO/WITHDRAW сообщения, маршрутизируемые по ETCP. * через NODEINFO/WITHDRAW сообщения, маршрутизируемые по ETCP.
@ -43,6 +49,7 @@ struct UTUN_INSTANCE;
struct NAT_DETECTION; struct NAT_DETECTION;
struct CONN_MGR; struct CONN_MGR;
struct TOPO_RECOVERY_CTX; struct TOPO_RECOVERY_CTX;
struct TOPO_GROUP_CONNECT;
/** Callback when a node's info (pubkeys, addresses) is persisted in the DB */ /** Callback when a node's info (pubkeys, addresses) is persisted in the DB */
typedef void (*topo_node_updated_fn)(struct UTUN_INSTANCE* inst, uint64_t node_id, typedef void (*topo_node_updated_fn)(struct UTUN_INSTANCE* inst, uint64_t node_id,
@ -132,6 +139,7 @@ struct TOPO_GROUP {
char channel_id[64]; // channel_id для групп типа CHAT char channel_id[64]; // channel_id для групп типа CHAT
struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы
struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления
struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT)
}; };
/** /**

260
src/routing_layer/topo_group_connect.c

@ -0,0 +1,260 @@
/*
* topo_group_connect.c — авто-подключение к узлам группы при старте
*
* Phase 1: одновременный запуск conn_mgr_connect_node для пиров с connected=1.
* Таймаут = direct_timeout_ms. По таймауту отменяем незавершённые и помечаем
* connected=0 в БД. Если хоть один подключился → done.
*
* Phase 2: последовательный перебор пиров с публичными адресами до первого успеха.
*
* Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями.
* UP: инкремент только для первого соединения к peer_node_id.
* DOWN: декремент только если нет других conn к peer_node_id И нет indirect-путей в nodeinfo.
*
* restart: перезапуск авто-подключения при phase=done и active_conn_count==0.
*/
#include "topo_group_connect.h"
#include "topo_group.h"
#include "topo_node.h"
#include "topo_node_sqlite.h"
#include "conn_mgr.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "etcp.h"
/* ─── внутренние константы ─── */
#define TGC_ID "topo_group_connect"
#define TGC_PHASE_ONE 0
#define TGC_PHASE_TWO 1
#define TGC_PHASE_DONE 2
struct TOPO_GROUP_CONNECT {
struct TOPO_GROUP* group;
void* phase_timer;
uint8_t phase;
uint8_t active;
int pending;
int connected_count;
int active_conn_count; // число уникальных peer-узлов с живыми соединениями
int cursor;
uint64_t* candidate_ids;
int candidate_count;
};
static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc);
static void tgc_phase2_result(int result, uint64_t node_id, void* arg);
static void tgc_phase1_timeout(void* arg);
static void tgc_phase1_result(int result, uint64_t node_id, void* arg);
/* ═══════════════════════════════════════════════════════════════════════
* Жизненный цикл
* ══════════════════════════════════════════════════════════════════════ */
int topo_group_connect_init(struct TOPO_GROUP* group) {
if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || group->connect)
return -1;
struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc));
if (!gc) return -1;
gc->group = group; gc->active = 1;
group->connect = gc;
uint64_t* ids = NULL; int count = 0;
sqlite3* db = group->instance->topo_sqlite_db;
if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) {
if (ids) { u_free(ids); ids = NULL; }
gc->phase = TGC_PHASE_TWO;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s no connected peers → Phase 2", TGC_ID, group->channel_id);
tgc_phase2_try_next(gc);
return 0;
}
gc->phase = TGC_PHASE_ONE; gc->pending = count;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s Phase 1 launching %d connects", TGC_ID, group->channel_id, count);
for (int i = 0; i < count; i++)
conn_mgr_connect_node(group->conn_mgr, ids[i], 0, tgc_phase1_result, gc);
u_free(ids);
gc->phase_timer = uasync_set_timeout(group->instance->ua,
(int)group->conn_mgr->direct_timeout_ms * 10, gc, tgc_phase1_timeout, "tgc_phase1");
return 0;
}
void topo_group_connect_destroy(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) return;
group->connect = NULL;
if (gc->phase_timer) { uasync_cancel_timeout(group->instance->ua, gc->phase_timer); gc->phase_timer = NULL; }
u_free(gc->candidate_ids);
u_free(gc);
}
void topo_group_connect_restart(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) { topo_group_connect_init(group); return; }
if (gc->phase != TGC_PHASE_DONE) return;
if (gc->active_conn_count > 0) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s", TGC_ID, group->channel_id);
topo_group_connect_destroy(group);
topo_group_connect_init(group);
}
int topo_group_connect_active_count(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect;
return gc ? gc->active_conn_count : 0;
}
/* ═══════════════════════════════════════════════════════════════════════
* ON UP / DOWN
* ══════════════════════════════════════════════════════════════════════ */
void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return;
uint64_t peer = conn->peer_node_id;
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) return;
/* уже есть другое прямое соединение к этому peer? */
struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx already connected via other conn, skip count",
TGC_ID, (unsigned long long)peer);
return;
}
e = e->next;
}
gc->active_conn_count++;
sqlite3* db = group->instance->topo_sqlite_db;
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s active=%d",
TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count);
}
void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return;
uint64_t peer = conn->peer_node_id;
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) return;
/* есть ли другие прямые соединения к этому peer? */
struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx still has other direct conn, skip",
TGC_ID, (unsigned long long)peer);
return;
}
e = e->next;
}
/* есть ли indirect-пути в nodeinfo? */
struct TOPO_NODEQ* nq = topo_node_find_by_id(group, peer);
if (nq && nq->paths) {
struct ll_entry* pe = nq->paths->head;
while (pe) {
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe;
if (path->conn && path->conn != conn) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx has indirect path, skip",
TGC_ID, (unsigned long long)peer);
return;
}
pe = pe->next;
}
}
gc->active_conn_count--;
sqlite3* db = group->instance->topo_sqlite_db;
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d",
TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count);
}
/* ═══════════════════════════════════════════════════════════════════════
* Phase 1
* ══════════════════════════════════════════════════════════════════════ */
static void tgc_phase1_result(int result, uint64_t node_id, void* arg) {
struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg;
if (gc->phase != TGC_PHASE_ONE) return;
gc->pending--;
if (result == CONN_MGR_OK) gc->connected_count++;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 result node=0x%016llx %s pending=%d connected=%d",
TGC_ID, (unsigned long long)node_id, result == CONN_MGR_OK ? "OK" : "FAIL", gc->pending, gc->connected_count);
}
static void tgc_phase1_timeout(void* arg) {
struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg;
gc->phase_timer = NULL;
if (gc->phase != TGC_PHASE_ONE) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 timeout connected=%d pending=%d ch=%s",
TGC_ID, gc->connected_count, gc->pending, gc->group->channel_id);
sqlite3* db = gc->group->instance->topo_sqlite_db;
uint64_t* ids = NULL; int count = 0;
topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count);
for (int i = 0; i < count; i++) {
uint8_t status; uint8_t conn_type;
conn_mgr_get_status(gc->group->conn_mgr, ids[i], &status, &conn_type);
if (status != CONN_MGR_STATE_CONNECTED) {
conn_mgr_cancel_callback(gc->group->conn_mgr, ids[i], tgc_phase1_result, gc);
topo_node_sqlite_set_connected(db, gc->group->channel_id, ids[i], 0);
}
}
u_free(ids);
if (gc->connected_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — %d connected, skipping Phase 2", TGC_ID, gc->connected_count);
gc->phase = TGC_PHASE_DONE;
return;
}
gc->phase = TGC_PHASE_TWO;
tgc_phase2_try_next(gc);
}
/* ═══════════════════════════════════════════════════════════════════════
* Phase 2
* ══════════════════════════════════════════════════════════════════════ */
static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) {
if (gc->phase != TGC_PHASE_TWO) return;
if (gc->candidate_count == 0) {
sqlite3* db = gc->group->instance->topo_sqlite_db;
topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count);
gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 loaded %d public peers for ch=%s",
TGC_ID, gc->candidate_count, gc->group->channel_id);
}
while (gc->cursor < gc->candidate_count) {
uint64_t nid = gc->candidate_ids[gc->cursor++];
uint8_t status; uint8_t conn_type;
conn_mgr_get_status(gc->group->conn_mgr, nid, &status, &conn_type);
if (status == CONN_MGR_STATE_CONNECTED) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 skip 0x%016llx already connected", TGC_ID, (unsigned long long)nid);
continue;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d)", TGC_ID,
(unsigned long long)nid, gc->cursor, gc->candidate_count);
conn_mgr_connect_node(gc->group->conn_mgr, nid, 0, tgc_phase2_result, gc);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id);
gc->phase = TGC_PHASE_DONE;
}
static void tgc_phase2_result(int result, uint64_t node_id, void* arg) {
struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg;
if (gc->phase != TGC_PHASE_TWO) return;
if (result == CONN_MGR_OK) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — done", TGC_ID, (unsigned long long)node_id);
gc->phase = TGC_PHASE_DONE;
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — try next", TGC_ID, (unsigned long long)node_id);
tgc_phase2_try_next(gc);
}

30
src/routing_layer/topo_group_connect.h

@ -0,0 +1,30 @@
/**
* @file topo_group_connect.h
* @brief Авто-подключение к узлам группы при старте и отслеживание состояния connected в БД.
*
* При старте CHAT-группы:
* Phase 1 — одновременный запуск подключений ко всем пирам с connected=1.
* Если хоть один подключился → done.
* Phase 2 — последовательный перебор по пирам с публичными адресами.
* До первого успешного подключения.
*
* Отслеживание: на UP → peers_<ch>.connected=1, на DOWN (последний conn) → 0.
* active_conn_count — количество уникальных peer-узлов с живыми соединениями.
* restart — перезапуск авто-подключения (если phase=done и нет подключений).
*/
#ifndef TOPO_GROUP_CONNECT_H
#define TOPO_GROUP_CONNECT_H
#include <stdint.h>
struct TOPO_GROUP;
struct ETCP_CONN;
int topo_group_connect_init(struct TOPO_GROUP* group);
void topo_group_connect_destroy(struct TOPO_GROUP* group);
void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
int topo_group_connect_active_count(struct TOPO_GROUP* group);
void topo_group_connect_restart(struct TOPO_GROUP* group);
#endif

63
src/routing_layer/topo_node_sqlite.c

@ -92,6 +92,9 @@ int topo_node_sqlite_init(sqlite3* db) {
snprintf(isql, sizeof(isql), snprintf(isql, sizeof(isql),
"CREATE INDEX IF NOT EXISTS \"idx_%s_type_rtt\" ON \"%s\" (node_type, node_RTT)", tbl, tbl); "CREATE INDEX IF NOT EXISTS \"idx_%s_type_rtt\" ON \"%s\" (node_type, node_RTT)", tbl, tbl);
sqlite3_exec(db, isql, NULL, NULL, NULL); sqlite3_exec(db, isql, NULL, NULL, NULL);
snprintf(isql, sizeof(isql),
"ALTER TABLE \"%s\" ADD COLUMN connected INTEGER NOT NULL DEFAULT 0", tbl);
sqlite3_exec(db, isql, NULL, NULL, NULL);
} }
sqlite3_finalize(s); sqlite3_finalize(s);
} }
@ -274,7 +277,8 @@ int topo_node_sqlite_channel_put(sqlite3* db, const char* channel_id,
" node_type INTEGER NOT NULL DEFAULT 0," /* вычисляемый тип узла (topo_node_sqlite_nodeinfo_updated): " node_type INTEGER NOT NULL DEFAULT 0," /* вычисляемый тип узла (topo_node_sqlite_nodeinfo_updated):
0=неизвестно 1=прямой(addr_type NETIF/DIRECT/NAT_STRICT) 0=неизвестно 1=прямой(addr_type NETIF/DIRECT/NAT_STRICT)
2=EIM NAT 3=(не исп.) 4=суперузел(DIRECT+"supernode=yes") */ 2=EIM NAT 3=(не исп.) 4=суперузел(DIRECT+"supernode=yes") */
" node_RTT INTEGER" /* RTT до узла в мс, часть индекса (node_type, node_RTT) */ " node_RTT INTEGER," /* RTT до узла в мс, часть индекса (node_type, node_RTT) */
" connected INTEGER NOT NULL DEFAULT 0" /* 1 = было активное подключение к узлу */
")"; ")";
int ddl_sz = snprintf(NULL, 0, ddl_tmpl, peers_tbl); int ddl_sz = snprintf(NULL, 0, ddl_tmpl, peers_tbl);
@ -724,3 +728,60 @@ struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* gr
(unsigned long long)node_id, name ? name : "", topo_list_count((struct _topo_head*)v4_head)); (unsigned long long)node_id, name ? name : "", topo_list_count((struct _topo_head*)v4_head));
return ni; return ni;
} }
int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected) {
if (!db || !channel_id) return -1;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[128]; snprintf(sql, sizeof(sql), "UPDATE \"%s\" SET connected=? WHERE node_id=?", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
sqlite3_bind_int(st, 1, connected ? 1 : 0);
sqlite3_bind_int64(st, 2, (sqlite3_int64)node_id);
int rc = sqlite3_step(st); sqlite3_finalize(st);
return (rc == SQLITE_DONE) ? 0 : -1;
}
int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id,
uint64_t** out_ids, int* out_count) {
if (!db || !channel_id || !out_ids || !out_count) return -1;
*out_ids = NULL; *out_count = 0;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[128]; snprintf(sql, sizeof(sql), "SELECT node_id FROM \"%s\" WHERE connected=1", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
int cnt = 0;
while (sqlite3_step(st) == SQLITE_ROW) cnt++;
if (cnt == 0) { sqlite3_finalize(st); return 0; }
uint64_t* ids = u_malloc((size_t)cnt * sizeof(uint64_t));
if (!ids) { sqlite3_finalize(st); return -1; }
sqlite3_reset(st); int i = 0;
while (sqlite3_step(st) == SQLITE_ROW) ids[i++] = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
*out_ids = ids; *out_count = cnt;
return 0;
}
int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id,
uint64_t** out_ids, int* out_count) {
if (!db || !channel_id || !out_ids || !out_count) return -1;
*out_ids = NULL; *out_count = 0;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[280];
snprintf(sql, sizeof(sql),
"SELECT p.node_id FROM \"%s\" p"
" WHERE EXISTS (SELECT 1 FROM node_addresses a"
" WHERE a.node_id=p.node_id AND a.family=4 AND a.addr_type!=0)"
" ORDER BY p.node_RTT", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
int cnt = 0;
while (sqlite3_step(st) == SQLITE_ROW) cnt++;
if (cnt == 0) { sqlite3_finalize(st); return 0; }
uint64_t* ids = u_malloc((size_t)cnt * sizeof(uint64_t));
if (!ids) { sqlite3_finalize(st); return -1; }
sqlite3_reset(st); int i = 0;
while (sqlite3_step(st) == SQLITE_ROW) ids[i++] = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
*out_ids = ids; *out_count = cnt;
return 0;
}

4
src/routing_layer/topo_node_sqlite.h

@ -67,4 +67,8 @@ void topo_node_sqlite_update_rtt(sqlite3* db, uint64_t node_id, uint16_t rtt);
*/ */
struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* groups, uint64_t node_id); struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* groups, uint64_t node_id);
int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected);
int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count);
int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count);
#endif #endif

3
tools/chatgui/CMakeLists.txt

@ -80,7 +80,6 @@ add_executable(chatgui
transport/config_updater.cpp transport/config_updater.cpp
transport/gui_bridge_impl.cpp transport/gui_bridge_impl.cpp
transport/chat_core.c transport/chat_core.c
transport/chat_conn_mgr.c
transport/chat_sync.c transport/chat_sync.c
transport/member_sync.c transport/member_sync.c
transport/merkle_sync.c transport/merkle_sync.c
@ -92,7 +91,7 @@ add_executable(chatgui
target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db) target_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) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE)
set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_conn_mgr.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/miniaudio_impl.c PROPERTIES LANGUAGE C) set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/miniaudio_impl.c PROPERTIES LANGUAGE C)
if(WIN32) if(WIN32)
target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread) target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread)
else() else()

4
tools/chatgui/libutun/CMakeLists.txt

@ -67,9 +67,7 @@ set(UTUN_COMMON_SOURCES
${SRC_DIR}/db_sync.c ${SRC_DIR}/db_sync.c
${SRC_DIR}/routing_layer/conn_mgr.c ${SRC_DIR}/routing_layer/conn_mgr.c
${SRC_DIR}/routing_layer/topo_recovery.c ${SRC_DIR}/routing_layer/topo_recovery.c
${TRANSPORT_DIR}/chat_sync.c ${SRC_DIR}/routing_layer/topo_group_connect.c
${TRANSPORT_DIR}/merkle_sync.c
${TRANSPORT_DIR}/member_sync.c
${SRC_DIR}/routing_layer/routing.c ${SRC_DIR}/routing_layer/routing.c
${SRC_DIR}/tun_if.c ${SRC_DIR}/tun_if.c
${SRC_DIR}/tun_route.c ${SRC_DIR}/tun_route.c

16
tools/chatgui/transport/chat_core.c

@ -6,11 +6,11 @@
*/ */
#include "chat_core.h" #include "chat_core.h"
#include "chat_conn_mgr.h"
#include "db_sync.h" #include "db_sync.h"
#include "gui_bridge.h" #include "gui_bridge.h"
#include "../../../lib/json_flat.h" #include "../../../lib/json_flat.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "topo_group.h"
#include "member_sync.h" #include "member_sync.h"
#include "../../../src/utun_instance.h" #include "../../../src/utun_instance.h"
@ -39,9 +39,8 @@
static struct chat_core_ctx { static struct chat_core_ctx {
struct UTUN_INSTANCE* inst; struct UTUN_INSTANCE* inst;
sqlite3* db; sqlite3* db;
uint8_t shared_db; /* db == inst->topo_sqlite_db, do not close */ uint8_t shared_db;
uint64_t my_node_id; uint64_t my_node_id;
struct chat_conn_mgr* conn_mgr;
/* db_sync instances per channel */ /* db_sync instances per channel */
struct DB_SYNC_INSTANCE** si; struct DB_SYNC_INSTANCE** si;
@ -136,7 +135,6 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
} }
} }
g_cc.conn_mgr = chat_conn_mgr_init(inst, g_cc.db);
g_cc.initialized = 1; g_cc.initialized = 1;
chat_core_sync_my_addresses(); chat_core_sync_my_addresses();
@ -237,16 +235,11 @@ int chat_core_is_initialized(void) {
return g_cc.initialized; 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 chat_core_destroy(struct UTUN_INSTANCE* inst) {
(void)inst; (void)inst;
if (!g_cc.initialized) return; if (!g_cc.initialized) return;
g_cc.initialized = 0; 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); } 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); DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID);
@ -670,6 +663,11 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
uint64_t ch_hash = 0; uint64_t ch_hash = 0;
{ const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); } { const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); }
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash); struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash);
/* create TOPO_GROUP_TYPE_CHAT for auto-connect */
uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid))
topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
if (si) { if (si) {
si_register(si, ch_id); si_register(si, ch_id);
db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id)); db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id));

2
tools/chatgui/transport/chat_core.h

@ -23,8 +23,6 @@ struct sqlite3* chat_core_get_db(void);
struct UTUN_INSTANCE* chat_core_get_inst(void); struct UTUN_INSTANCE* chat_core_get_inst(void);
int chat_core_is_initialized(void); int chat_core_is_initialized(void);
struct chat_conn_mgr* chat_conn_mgr_get(void);
/* ── Отправка сообщения (GUI → uasync) ── */ /* ── Отправка сообщения (GUI → uasync) ── */
struct chat_msg_submit { struct chat_msg_submit {

333
tools/chatgui/transport/chat_sync.c

@ -1,8 +1,10 @@
#include "chat_sync.h" #include "chat_sync.h"
#include "chat_core.h" #include "chat_core.h"
#include "chat_conn_mgr.h"
#include "gui_bridge.h" #include "gui_bridge.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "topo_group.h"
#include "conn_mgr.h"
#include "topo_group_connect.h"
#include "../../../lib/json_flat.h" #include "../../../lib/json_flat.h"
#include "member_sync.h" #include "member_sync.h"
#include "merkle_sync.h" #include "merkle_sync.h"
@ -24,238 +26,6 @@
static struct chat_sync* g_cs = NULL; static struct chat_sync* g_cs = NULL;
/* ═══════════════════════════════════════════════════════════════════════
* Auto-connect: cursor-based peer iteration with GC every 1s
* ══════════════════════════════════════════════════════════════════════ */
#define AC_MAX_FLIGHTS 10
#define AC_RETRY_MS 1000
#define AC_GC_TIMEOUT_MS 3000
#define AC_ID "auto_connect"
struct ac_flight {
uint64_t node_id;
uint8_t active; /* 1 = connection attempt in progress */
uint64_t created_tb;
};
struct auto_connect {
struct UTUN_INSTANCE* inst;
void* retry_timer;
uint8_t active;
char** channel_ids; /* loaded once, refreshed on cursor wrap */
int channel_count;
int ch_cursor; /* current channel index */
int peer_cursor; /* current peer index within channel */
struct ac_flight flights[AC_MAX_FLIGHTS];
};
static struct auto_connect* g_ac = NULL;
static void ac_retry_timer_cb(void* arg);
static void ac_result_cb(int result, uint64_t node_id, void* arg);
/* ── load channel IDs into auto_connect ── */
static int ac_load_channels(struct auto_connect* ac) {
if (ac->channel_ids) {
for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]);
u_free(ac->channel_ids);
ac->channel_ids = NULL;
}
ac->channel_count = 0;
uint8_t buf[4096]; size_t buf_len;
if (chat_core_list_channels(buf, sizeof(buf), &buf_len) != 0 || buf_len < 2) return 0;
uint16_t cnt; memcpy(&cnt, buf, 2);
const uint8_t* p = buf + 2; size_t rem = buf_len - 2;
ac->channel_ids = u_calloc(cnt, sizeof(char*));
if (!ac->channel_ids) return 0;
for (uint16_t i = 0; i < cnt && rem >= 1; i++) {
uint8_t id_len = *p++; rem--;
if (rem < id_len) break;
ac->channel_ids[i] = u_malloc(id_len + 1);
if (ac->channel_ids[i]) { memcpy(ac->channel_ids[i], p, id_len); ac->channel_ids[i][id_len] = '\0'; ac->channel_count++; }
p += id_len; rem -= id_len;
}
return ac->channel_count;
}
/* ── GC: close expired flights (no link_status after AC_GC_TIMEOUT_MS) ── */
static void ac_gc(struct auto_connect* ac) {
uint64_t now = get_time_tb();
uint64_t deadline = (uint64_t)AC_GC_TIMEOUT_MS * 10;
for (int i = 0; i < AC_MAX_FLIGHTS; i++) {
if (!ac->flights[i].active) continue;
if (now - ac->flights[i].created_tb < deadline) continue;
uint64_t nid = ac->flights[i].node_id;
/* check if link already UP — if so, just free slot (already connected) */
int found_up = 0;
{
struct ll_entry* entry = ac->inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (ce->conn->peer_node_id == nid) {
struct ETCP_LINK* l = ce->conn->links;
while (l) { if (l->initialized && l->link_status) { found_up = 1; break; } l = l->next; }
if (found_up) break;
}
entry = entry->next;
}
}
if (found_up) {
chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id);
memset(&ac->flights[i], 0, sizeof(ac->flights[i]));
continue;
}
/* expired and not up — cancel */
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: GC closing flight node=0x%016llx (expired)", AC_ID, (unsigned long long)nid);
chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id);
memset(&ac->flights[i], 0, sizeof(ac->flights[i]));
}
}
/* ── count occupied flight slots ── */
static int ac_flight_count(struct auto_connect* ac) {
int n = 0;
for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (ac->flights[i].active) n++;
return n;
}
/* ── find first free flight slot ── */
static int ac_find_free_slot(struct auto_connect* ac) {
for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (!ac->flights[i].active) return i;
return -1;
}
/* ── check if node already has active ETCP connection ── */
static int ac_node_has_conn(struct UTUN_INSTANCE* inst, uint64_t nid) {
if (!inst->connections) return 0;
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&nid);
return e != NULL;
}
/* ── fill up to AC_MAX_FLIGHTS by advancing cursor ── */
static void ac_fill(struct auto_connect* ac) {
if (ac->channel_count == 0) return;
/* try all channels up to 2 full cycles without finding a candidate → give up */
int max_iter = ac->channel_count * 2 + 10;
while (ac_flight_count(ac) < AC_MAX_FLIGHTS && max_iter-- > 0) {
int slot = ac_find_free_slot(ac);
if (slot < 0) break;
char* ch_id = ac->channel_ids[ac->ch_cursor];
uint8_t pb[2048]; size_t plen;
if (chat_core_list_peers(ch_id, pb, sizeof(pb), &plen) != 0 || plen < 2) {
/* channel has no peers or error — skip to next */
ac->ch_cursor++; ac->peer_cursor = 0;
if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); }
continue;
}
uint16_t pc; memcpy(&pc, pb, 2);
if (ac->peer_cursor >= (int)pc) {
ac->ch_cursor++; ac->peer_cursor = 0;
if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); }
continue;
}
/* get next peer */
uint64_t nid; memcpy(&nid, pb + 2 + ac->peer_cursor * 8, 8);
ac->peer_cursor++;
if (nid == 0 || nid == ac->inst->node_id) continue;
/* already connected? */
if (ac_node_has_conn(ac->inst, nid)) continue;
/* already in a flight slot? */
int dup = 0;
for (int i = 0; i < AC_MAX_FLIGHTS; i++)
if (ac->flights[i].active && ac->flights[i].node_id == nid) { dup = 1; break; }
if (dup) continue;
/* launch */
struct ac_flight* f = &ac->flights[slot];
f->node_id = nid;
f->created_tb = get_time_tb();
chat_conn_mgr_connect(chat_conn_mgr_get(), nid, ac_result_cb, f); f->active = 1;
if (f->node_id != nid) {
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: flight to 0x%016llx FAIL (ca_state=NULL, sync fail in chat_conn_mgr_connect_auto)", AC_ID, (unsigned long long)nid);
f->node_id = 0; f->created_tb = 0;
continue;
}
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: flight[%d] launched node=0x%016llx", AC_ID, slot, (unsigned long long)nid);
uint8_t evt[7]; evt[0] = 0; uint16_t t = 0, n = (uint16_t)ac->channel_count, s = 0;
memcpy(evt + 1, &t, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &s, 2);
gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7);
}
}
/* ── result callback (fired by chat_conn_mgr_connect_auto) ── */
static void ac_result_cb(int result, uint64_t node_id, void* arg) {
struct ac_flight* f = (struct ac_flight*)arg;
if (!g_ac || !g_ac->active) return;
const char* rs = result == CC_OK ? "OK" : "FAIL";
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: result node=0x%016llx %s", AC_ID, (unsigned long long)node_id, rs);
f->active = 0; f->node_id = 0; f->created_tb = 0;
}
/* ── retry timer callback (every AC_RETRY_MS) ── */
static void ac_retry_timer_cb(void* arg) {
struct auto_connect* ac = (struct auto_connect*)arg;
ac->retry_timer = NULL;
if (!ac->active) return;
if (ac->channel_count == 0) ac_load_channels(ac);
ac_gc(ac);
ac_fill(ac);
ac->retry_timer = uasync_set_timeout(ac->inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry");
}
/* ═══════════════════════════════════════════════════════════════════════
* Public auto_connect API
* ══════════════════════════════════════════════════════════════════════ */
void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->ua) return;
chat_sync_auto_connect_stop();
struct auto_connect* ac = u_calloc(1, sizeof(*ac));
if (!ac) return;
ac->inst = inst;
ac->active = 1;
g_ac = ac;
ac_load_channels(ac);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: start, %d channels loaded", AC_ID, ac->channel_count);
ac_gc(ac);
ac_fill(ac);
ac->retry_timer = uasync_set_timeout(inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry");
}
void chat_sync_auto_connect_stop(void) {
struct auto_connect* ac = g_ac;
if (!ac) return;
ac->active = 0; g_ac = NULL;
if (ac->retry_timer) { uasync_cancel_timeout(ac->inst->ua, ac->retry_timer); ac->retry_timer = NULL; }
for (int i = 0; i < AC_MAX_FLIGHTS; i++)
if (ac->flights[i].active) { chat_conn_mgr_cancel(chat_conn_mgr_get(), ac->flights[i].node_id); memset(&ac->flights[i], 0, sizeof(ac->flights[i])); }
if (ac->channel_ids) { for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]); u_free(ac->channel_ids); }
uint8_t evt[7]; evt[0] = 2; uint16_t z = 0;
memcpy(evt + 1, &z, 2); memcpy(evt + 3, &z, 2); memcpy(evt + 5, &z, 2);
gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7);
u_free(ac);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: stopped", AC_ID);
}
void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: switch group 0x%016llx", AC_ID, (unsigned long long)new_group_id);
chat_sync_auto_connect_stop();
chat_sync_auto_connect_start(inst);
}
struct channel_cache { struct channel_cache {
char channel_id[64]; char channel_id[64];
uint64_t* peer_ids; uint64_t* peer_ids;
@ -399,7 +169,7 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
case CS_MSG_ERROR: { case CS_MSG_ERROR: {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id); DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id);
if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; } if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; }
uint8_t err[20]; memcpy(err, &peer, 8); int r = CC_ERR_NOT_FOUND; memcpy(err + 8, &r, 4); uint8_t err[20]; memcpy(err, &peer, 8); int r = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4);
memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); 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; g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
break; break;
@ -420,7 +190,7 @@ static void cs_info_req_timeout_cb(void* arg) {
CS_ID, (unsigned long long)cs->pending_invite_node_id, CS_ID, (unsigned long long)cs->pending_invite_node_id,
(unsigned long long)cs->pending_invite_ch_id); (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8); uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4); int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
@ -435,7 +205,7 @@ static void cs_join_timeout_cb(void* arg) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu", DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id); CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &peer, 8); uint8_t err[20]; memcpy(err, &peer, 8);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4); int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
@ -785,7 +555,6 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
struct chat_sync* cs = g_cs; struct chat_sync* cs = g_cs;
if (!cs || !inst) return; if (!cs || !inst) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: destroy", CS_ID); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: destroy", CS_ID);
chat_sync_auto_connect_stop();
if (cs->sync_timer) { uasync_cancel_timeout(inst->ua, cs->sync_timer); cs->sync_timer = NULL; } if (cs->sync_timer) { uasync_cancel_timeout(inst->ua, cs->sync_timer); cs->sync_timer = NULL; }
if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); } if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); }
member_sync_destroy(inst); member_sync_destroy(inst);
@ -803,15 +572,87 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
} }
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
(void)inst; (void)inst; (void)node_id;
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; }; struct chat_invite {
uint64_t channel_id;
uint64_t node_id;
uint8_t pubkey[32];
uint8_t* addrs_data;
int addr_count;
};
struct cm_invite_wrap { struct chat_invite inv; };
static void cm_invite_trampoline(void* arg) { static void cm_invite_trampoline(void* arg) {
struct cm_invite_wrap* w = (struct cm_invite_wrap*)arg; struct cm_invite_wrap* w = (struct cm_invite_wrap*)arg;
chat_conn_mgr_connect_from_invite(w->mgr, &w->inv); struct chat_invite* inv = &w->inv;
struct UTUN_INSTANCE* inst = chat_core_get_inst();
sqlite3* db = chat_core_get_db();
uint64_t node_id = inv->node_id;
uint64_t channel_id = inv->channel_id;
/* save pubkey to nodes */
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db, "INSERT INTO nodes(node_id,x25519_pubkey,created_at) VALUES(?,?,?)"
" ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey,"
" created_at=COALESCE(nodes.created_at, excluded.created_at)", -1, &st, NULL);
if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 3, 0); sqlite3_step(st); sqlite3_finalize(st); }
/* save invite addresses (socket_id=0) */
sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", -1, &st, NULL);
if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_step(st); sqlite3_finalize(st); }
sqlite3_prepare_v2(db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id) VALUES(?,?,1,?,?,0,?)", -1, &st, NULL);
if (st) {
const uint8_t* ap = inv->addrs_data;
for (int i = 0; i < inv->addr_count; i++) {
uint8_t family = *ap++; ap++; /* skip sock_id */
if (family != 4) { ap += 18; continue; }
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_int(st, 2, 4);
sqlite3_bind_blob(st, 3, ap, 4, SQLITE_STATIC); ap += 4;
uint16_t port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2;
sqlite3_bind_int(st, 4, (int)port); sqlite3_bind_int(st, 5, 1); /* addr_type=1 (DIRECT) */
sqlite3_step(st); sqlite3_reset(st);
}
sqlite3_finalize(st);
}
char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)channel_id);
chat_core_ensure_channel_ready(ch_str);
/* ensure TOPO_GROUP_TYPE_CHAT group exists */
if (inst->topo_groups) {
uint64_t gid = channel_id;
if (!topo_groups_find(inst->topo_groups, gid))
topo_groups_create_group(inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_str);
}
/* build TOPO_NODE from invite data */
struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni));
if (!ni) { u_free(w); return; }
ni->node_id = node_id; memcpy(ni->public_key, inv->pubkey, 32);
const uint8_t* ap = inv->addrs_data;
for (int i = 0; i < inv->addr_count; i++) {
uint8_t family = *ap++; ap++; /* skip sock_id */
if (family != 4) { ap += 18; continue; }
struct TOPO_ADDR4* a4 = u_calloc(1, sizeof(*a4));
if (!a4) break;
memcpy(a4->addr, ap, 4); ap += 4;
a4->port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2;
a4->type = TOPO_ADDR_NAT; a4->protocol = TOPO_PROTO_UDP;
a4->next = ni->v4_addrs; ni->v4_addrs = a4;
}
if (inst->topo_groups && ni->v4_addrs) {
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id);
if (g && g->conn_mgr)
conn_mgr_connect_from_invite(g->conn_mgr, ni, channel_id, 0, NULL, NULL);
}
/* free caller-owned ni */
while (ni->v4_addrs) { struct TOPO_ADDR4* n = ni->v4_addrs->next; u_free(ni->v4_addrs); ni->v4_addrs = n; }
u_free(ni);
u_free(w); u_free(w);
} }
@ -846,8 +687,8 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start ch=%llu node=0x%016llx pubkey=%016llx... addrs=%d", 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); 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)); struct cm_invite_wrap { struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap));
w->mgr = chat_conn_mgr_get(); w->inv = *inv; u_free(inv); w->inv = *inv; u_free(inv);
gui_bridge_post_uasync_fn( gui_bridge_post_uasync_fn(
(void(*)(void*))cm_invite_trampoline, w); (void(*)(void*))cm_invite_trampoline, w);
} }
@ -1193,7 +1034,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(joiner) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc); if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(joiner) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
{ char j_name[128]; json_flat_get(joiner_userinfo, "name", j_name, sizeof(j_name)); { char j_name[128]; json_flat_get(joiner_userinfo, "name", j_name, sizeof(j_name));
topo_node_sqlite_node_update_verified(db, node_id, j_name, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } topo_node_sqlite_node_update_verified(db, node_id, j_name, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); }
chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id);
/* save/update node addresses */ /* save/update node addresses */
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN save addrs node=0x%016llx addr_cnt=%d", CS_ID, (unsigned long long)node_id, addr_cnt); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN save addrs node=0x%016llx addr_cnt=%d", CS_ID, (unsigned long long)node_id, addr_cnt);
@ -1348,7 +1189,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
{ char pn[256]; json_flat_get(peer_userinfo, "name", pn, sizeof(pn)); { char pn[256]; json_flat_get(peer_userinfo, "name", pn, sizeof(pn));
topo_node_sqlite_node_update_verified(db, node_id, pn, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } topo_node_sqlite_node_update_verified(db, node_id, pn, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); }
chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] WELCOME save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); 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++) { for (uint8_t j = 0; j < ac; j++) {
@ -1451,7 +1292,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(peer_upsert) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc); if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(peer_upsert) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
{ char pn2[128]; json_flat_get(peer_userinfo, "name", pn2, sizeof(pn2)); { char pn2[128]; json_flat_get(peer_userinfo, "name", pn2, sizeof(pn2));
topo_node_sqlite_node_update_verified(db, node_id, pn2, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } topo_node_sqlite_node_update_verified(db, node_id, pn2, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); }
chat_conn_mgr_add_node(chat_conn_mgr_get(), node_id);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac);
for (uint8_t i = 0; i < ac; i++) { for (uint8_t i = 0; i < ac; i++) {

5
tools/chatgui/transport/chat_sync.h

@ -54,11 +54,6 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
const uint8_t* pubkey_bin, const uint8_t* pubkey_bin,
const uint8_t* addrs_data, int addr_count); const uint8_t* addrs_data, int addr_count);
/* Авто-подключение к узлам группы (параллельно до 10, останавливается при 3 успешных) */
void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst);
void chat_sync_auto_connect_stop(void);
void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id);
#ifdef __cplusplus #ifdef __cplusplus
} }
#endif #endif

25
tools/chatgui/transport/utun_node.cpp

@ -28,6 +28,7 @@ extern "C" {
#include "gui_bridge.h" #include "gui_bridge.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "topo_node.h" #include "topo_node.h"
#include "topo_group.h"
} }
#define ETCP_RT_ID_CHAT 0x11 #define ETCP_RT_ID_CHAT 0x11
@ -252,7 +253,6 @@ void UtunNode::runLoop() {
} }
utun_instance_set_tun_init_enabled(0); utun_instance_set_tun_init_enabled(0);
utun_instance_set_topo_group_enabled(0);
struct UASYNC* ua = uasync_create(); struct UASYNC* ua = uasync_create();
if (!ua) { if (!ua) {
@ -315,8 +315,27 @@ void UtunNode::runLoop() {
/* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */ /* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */
chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData()); chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData());
chat_sync_init(m_instance, nullptr); chat_sync_init(m_instance, nullptr);
chat_sync_auto_connect_start(m_instance);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync + auto_connect initialized"); /* create TOPO_GROUP_TYPE_CHAT for all existing channels */
if (m_instance->topo_groups) {
uint8_t buf[4096]; size_t ch_len = 0;
chat_core_list_channels(buf, sizeof(buf), &ch_len);
if (ch_len >= 2) {
uint16_t cnt; memcpy(&cnt, buf, 2);
const uint8_t* p = buf + 2; size_t rem = ch_len - 2;
for (uint16_t i = 0; i < cnt && rem >= 1; i++) {
uint8_t id_len = *p++; rem--;
if (rem < id_len) break;
char ch_id[64]; memcpy(ch_id, p, id_len < 63 ? id_len : 63); ch_id[id_len < 63 ? id_len : 63] = 0;
p += id_len; rem -= id_len;
uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid != 0 && !topo_groups_find(m_instance->topo_groups, gid))
topo_groups_create_group(m_instance->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
}
}
}
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized");
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop");
while (!m_stop) { while (!m_stop) {

Loading…
Cancel
Save