Browse Source

Use shared group requests for config and chat connection planners

proxy
evgeny 3 days ago
parent
commit
5838e63cde
  1. 2
      src/chat/chat_sync.c
  2. 3
      src/chat/member_sync.c
  3. 57
      src/routing_layer/topo_group.c
  4. 7
      src/routing_layer/topo_group.h
  5. 846
      src/routing_layer/topo_group_connect.c
  6. 36
      src/routing_layer/topo_group_connect.h
  7. 49
      src/routing_layer/topo_group_invite.c
  8. 39
      src/transport_layer/etcp_connections.c
  9. 2
      src/utun_instance.c
  10. 2
      src/utun_instance.h
  11. 16
      tests/bbr_integration/test_bbr_integration.c
  12. 38
      tests/test_group_recovery.c
  13. 7
      tests/test_ncd_config.c

2
src/chat/chat_sync.c

@ -1440,7 +1440,7 @@ static void cs_handle_join_ready(struct chat_sync* cs, uint64_t peer,
struct TOPO_GROUP* g = topo_groups_find(cs->inst->topo_groups, group_id); struct TOPO_GROUP* g = topo_groups_find(cs->inst->topo_groups, group_id);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) { if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) {
struct ETCP_CONN* c = cs_find_conn_for_node(cs->inst, peer); struct ETCP_CONN* c = cs_find_conn_for_node(cs->inst, peer);
if (c) topo_group_new_conn(g, c); if (c) topo_group_connect_node_once(g, c->peer_node_id);
} }
} }
chat_core_ensure_channel_ready(cs->inst, ch_id); chat_core_ensure_channel_ready(cs->inst, ch_id);

3
src/chat/member_sync.c

@ -1,3 +1,4 @@
#include "../routing_layer/topo_group_connect.h"
#include "member_sync.h" #include "member_sync.h"
#include "../routing_layer/topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
@ -706,7 +707,7 @@ static void member_committed(struct UTUN_INSTANCE* inst, const char* ns, uint64_
merkle_sync_changed(inst, ns); merkle_sync_changed(inst, ns);
struct TOPO_GROUP* g = inst->topo_groups ? topo_groups_find(inst->topo_groups, strtoull(ns, NULL, 10)) : NULL; struct TOPO_GROUP* g = inst->topo_groups ? topo_groups_find(inst->topo_groups, strtoull(ns, NULL, 10)) : NULL;
struct ETCP_CONN* conn = instance_find_conn(inst, node_id); struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT && conn) topo_group_new_conn(g, conn); if (g && g->group_type == TOPO_GROUP_TYPE_CHAT && conn) topo_group_connect_node_once(g, node_id);
} }
int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id, int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id,

57
src/routing_layer/topo_group.c

@ -71,6 +71,7 @@ static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) {
} }
if (peer->handle) node_conn_direct_close(peer->handle); if (peer->handle) node_conn_direct_close(peer->handle);
topo_recovery_changed(peer->group); topo_recovery_changed(peer->group);
topo_group_connect_changed(peer->group);
} }
static void topo_group_pending_remove(struct TOPO_GROUP_CONN_ITEM* peer) { static void topo_group_pending_remove(struct TOPO_GROUP_CONN_ITEM* peer) {
@ -93,7 +94,10 @@ static void topo_group_peer_event(struct NODE_CONN_DIRECT* handle, enum ncd_even
DEBUG_WARN(DEBUG_CATEGORY_BGP, "group peer transport timeout peer=%016llx group=%016llx", DEBUG_WARN(DEBUG_CATEGORY_BGP, "group peer transport timeout peer=%016llx group=%016llx",
(unsigned long long)peer->node_id, (unsigned long long)peer->group->group_id); (unsigned long long)peer->node_id, (unsigned long long)peer->group->group_id);
if (peer->conn) topo_group_remove_conn(peer->group, peer->conn, TOPO_REMOVE_TRANSPORT_DOWN); if (peer->conn) topo_group_remove_conn(peer->group, peer->conn, TOPO_REMOVE_TRANSPORT_DOWN);
else topo_group_pending_remove(peer); else if ((event == NCD_EVENT_DOWN || event == NCD_EVENT_TIMEOUT) && peer->requests && node_conn_direct_get_conn(handle)) {
peer->transport_down = 1; peer->retained = 0; topo_recovery_changed(peer->group);
topo_group_connect_changed(peer->group);
} else topo_group_pending_remove(peer);
} }
int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out) { int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out) {
@ -125,11 +129,22 @@ int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO
enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request) { enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request) {
struct TOPO_GROUP_CONN_ITEM* peer = request ? request->peer : NULL; struct TOPO_GROUP_CONN_ITEM* peer = request ? request->peer : NULL;
if (!peer || peer->tx_failed) return TOPO_PEER_FAILED; if (!peer || peer->tx_failed || peer->transport_down) return TOPO_PEER_FAILED;
if (!peer->conn || !peer->conn->links_up) return TOPO_PEER_CONNECTING; if (!peer->conn || !peer->conn->links_up) return TOPO_PEER_CONNECTING;
return topo_group_peer_ready(peer->group, peer->node_id) ? TOPO_PEER_READY : TOPO_PEER_SYNCING; return topo_group_peer_ready(peer->group, peer->node_id) ? TOPO_PEER_READY : TOPO_PEER_SYNCING;
} }
struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request) {
return request && request->peer ? node_conn_direct_get_conn(request->peer->handle) : NULL;
}
void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id) {
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id);
if (!peer) return;
if (peer->conn) topo_group_leave_peer(group, peer->handle);
else topo_group_pending_remove(peer);
}
uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request) { uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request) {
return request && request->peer ? request->peer->progress : 0; return request && request->peer ? request->peer->progress : 0;
} }
@ -153,10 +168,13 @@ static void topo_group_peer_progressed(struct TOPO_GROUP* group, uint64_t node_i
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id); struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id);
if (peer) peer->progress = get_time_tb(); if (peer) peer->progress = get_time_tb();
topo_recovery_changed(group); topo_recovery_changed(group);
topo_group_connect_changed(group);
} }
static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_CONN_ITEM* peer) { static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_CONN_ITEM* peer) {
if (topo_group_peer_ready(group, peer->node_id)) peer->retained = 1; if (topo_group_peer_ready(group, peer->node_id)) {
peer->retained = 1; topo_group_connect_on_up(group, peer->conn);
}
topo_group_peer_progressed(group, peer->node_id); topo_group_peer_progressed(group, peer->node_id);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: group=%016llx peer=%016llx exchange=%016llx sent=%u received=%u ready=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: group=%016llx peer=%016llx exchange=%016llx sent=%u received=%u ready=%d",
(unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id, (unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id,
@ -192,7 +210,7 @@ static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t si
} }
static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) {
topo_tx_clear(peer); peer->tx_failed = 0; topo_tx_clear(peer); peer->tx_failed = 0; peer->transport_down = 0;
peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0; peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0;
uint64_t epoch; uint64_t epoch;
if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) { if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) {
@ -329,10 +347,10 @@ static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t nod
// Приём пакетов // Приём пакетов
// ============================================================================ // ============================================================================
static struct NODE_CONN_DIRECT* topo_config_handle(struct UTUN_INSTANCE* instance, uint64_t node_id) { static int topo_configured(struct UTUN_INSTANCE* instance, uint64_t node_id) {
for (struct CONFIG_CONN_HANDLE* owner = instance->config_conn_handles; owner; owner = owner->next) for (struct CONFIG_CONN_HANDLE* owner = instance->config_conn_handles; owner; owner = owner->next)
if (owner->node_id == node_id) return owner->handle; if (owner->node_id == node_id) return 1;
return NULL; return 0;
} }
static int topo_group_has_sender(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { static int topo_group_has_sender(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
@ -473,7 +491,7 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg
if (status == ETCP_CONN_STATUS_UP && conn->peer_node_id) { if (status == ETCP_CONN_STATUS_UP && conn->peer_node_id) {
/* UTUN присоединяет пир через UP принадлежащего ему NCD handle. /* UTUN присоединяет пир через UP принадлежащего ему NCD handle.
* Сам по себе транспортный UP не означает участия в UTUN-группе. */ * Сам по себе транспортный UP не означает участия в UTUN-группе. */
if (!topo_config_handle(conn->instance, conn->peer_node_id)) topo_group_send_resync(conn); if (!topo_configured(conn->instance, conn->peer_node_id)) topo_group_send_resync(conn);
/* CHAT: только члены каналов (peers_* в БД). */ /* CHAT: только члены каналов (peers_* в БД). */
sqlite3* db = conn->instance ? conn->instance->topo_sqlite_db : NULL; sqlite3* db = conn->instance ? conn->instance->topo_sqlite_db : NULL;
uint64_t* chs = NULL; int chn = 0; uint64_t* chs = NULL; int chn = 0;
@ -780,8 +798,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
} }
u_free(chs); u_free(chs);
} else if (group_type == TOPO_GROUP_TYPE_UTUN && cqe->conn->links_up) { } else if (group_type == TOPO_GROUP_TYPE_UTUN && cqe->conn->links_up) {
struct NODE_CONN_DIRECT* handle = topo_config_handle(g->instance, cqe->peer_node_id); if (topo_configured(g->instance, cqe->peer_node_id)) topo_group_new_conn(group, cqe->conn);
if (handle) topo_group_join_peer(group, handle);
} }
} }
entry = entry->hash_next; entry = entry->hash_next;
@ -909,7 +926,11 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en
if (pending && !pending->conn && node_conn_direct_get_conn(pending->handle) == conn) { if (pending && !pending->conn && node_conn_direct_get_conn(pending->handle) == conn) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group pending peer removed: peer=%016llx group=%016llx reason=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "group pending peer removed: peer=%016llx group=%016llx reason=%d",
(unsigned long long)pending->node_id, (unsigned long long)group->group_id, reason); (unsigned long long)pending->node_id, (unsigned long long)group->group_id, reason);
topo_group_pending_remove(pending); return; if (reason == TOPO_REMOVE_TRANSPORT_DOWN && pending->requests) {
pending->transport_down = 1; pending->retained = 0;
topo_recovery_changed(group); topo_group_connect_changed(group);
} else topo_group_pending_remove(pending);
return;
} }
bool found_in_list = false; bool found_in_list = false;
@ -983,9 +1004,13 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en
while (e) { while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn == conn) { if (item->conn == conn) {
queue_remove_data(group->senders_list, e); if (reason == TOPO_REMOVE_TRANSPORT_DOWN && item->requests && node_conn_direct_get_conn(item->handle) == conn) {
topo_group_peer_free(item); topo_tx_clear(item); item->conn = NULL; item->transport_down = 1; item->retained = 0;
queue_entry_free(e); item->accepted = 0; item->table_sent = 0; item->table_received = 0;
topo_recovery_changed(group); topo_group_connect_changed(group);
} else {
queue_remove_data(group->senders_list, e); topo_group_peer_free(item); queue_entry_free(e);
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group_remove_conn: REMOVED conn=%p peer=%016llx from senders_list grp=%016llx ch=%s", DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group_remove_conn: REMOVED conn=%p peer=%016llx from senders_list grp=%016llx ch=%s",
(void*)conn, (unsigned long long)conn->peer_node_id, (void*)conn, (unsigned long long)conn->peer_node_id,
(unsigned long long)group->group_id, group->channel_id); (unsigned long long)group->group_id, group->channel_id);
@ -1163,6 +1188,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets }; struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets };
topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected); topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected);
u_free(new_hop_list); u_free(new_hop_list);
topo_group_peer_progressed(group, from->peer_node_id);
return 0; return 0;
} }
uint64_t incoming_timestamp = new_ni->timestamp; uint64_t incoming_timestamp = new_ni->timestamp;
@ -1388,6 +1414,7 @@ static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer) {
(unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch); (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch);
topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0; topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0;
topo_recovery_changed(peer->group); topo_recovery_changed(peer->group);
topo_group_connect_changed(peer->group);
} }
static void topo_tx_wait(void* arg); static void topo_tx_wait(void* arg);
@ -1694,7 +1721,7 @@ static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_C
static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn) { static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn) {
if (!instance || !conn || !conn->instance || !instance->topo_groups) return; if (!instance || !conn || !conn->instance || !instance->topo_groups) return;
/* общий сигнал от пассивной стороны: рестартуем обмен, если мы для этого пира VPN-клиент */ /* общий сигнал от пассивной стороны: рестартуем обмен, если мы для этого пира VPN-клиент */
if (!topo_config_handle(instance, conn->peer_node_id)) return; if (!topo_configured(instance, conn->peer_node_id)) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_resync: from %s — re-announce groups", conn->log_name); DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_resync: from %s — re-announce groups", conn->log_name);
struct ll_entry* fe = instance->topo_groups->group_list->head; struct ll_entry* fe = instance->topo_groups->group_list->head;
while (fe) { while (fe) {

7
src/routing_layer/topo_group.h

@ -179,7 +179,7 @@ struct TOPO_GROUP_CONN_ITEM {
uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint8_t retained; // участие запрошено постоянным владельцем или достигло READY
uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения
uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch
uint8_t tx_failed; uint8_t tx_failed, transport_down;
struct topo_tx_item *tx_head, *tx_tail; struct topo_tx_item *tx_head, *tx_tail;
struct queue_waiter_handle tx_waiter; struct queue_waiter_handle tx_waiter;
void* tx_wake; void* tx_wake;
@ -306,11 +306,16 @@ enum topo_peer_phase { TOPO_PEER_FAILED, TOPO_PEER_CONNECTING, TOPO_PEER_SYNCING
* владеет транспортом. Проверка CHAT-членства общая с входящим JOIN_GROUP. * владеет транспортом. Проверка CHAT-членства общая с входящим JOIN_GROUP.
* close отменяет только свой запрос; последний запрос отменяет незавершённую * close отменяет только свой запрос; последний запрос отменяет незавершённую
* попытку, если её не удерживает другой инициатор. READY сохраняется в группе. * попытку, если её не удерживает другой инициатор. READY сохраняется в группе.
* DOWN/TIMEOUT сообщает FAILED, но живой запрос сохраняет интерес к NCD: поздний
* UP возобновляет сессию. Закрытие последнего запроса отменяет неудачную попытку.
* CLOSED/leave/удаление группы отсоединяет запрос окончательно.
* Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан * Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан
* закрыть его. Все операции выполняются в потоке единственного uasync. * закрыть его. Все операции выполняются в потоке единственного uasync.
* progress — timebase последнего принятого NODEINFO/шага обмена, не polling time. */ * progress — timebase последнего принятого NODEINFO/шага обмена, не polling time. */
int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out); int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out);
void topo_group_peer_close(struct TOPO_PEER_REQUEST* request); void topo_group_peer_close(struct TOPO_PEER_REQUEST* request);
struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request); /* borrowed, including CONNECTING */
void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id);
enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request); enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request);
uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request); uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request);

846
src/routing_layer/topo_group_connect.c

@ -1,679 +1,283 @@
/* /* Политика CHAT-подключений. Транспортом и JOIN владеет групповая сессия.
* topo_group_connect.c — авто-подключение к узлам группы (бесконечный цикл до цели) * Сначала параллельно восстанавливаем connected-пиров, затем последовательно
* * пробуем supernode/public/local. Незавершённые попытки — отменяемые запросы. */
* Phase 1: одновременный запуск node_conn_direct_open для пиров с connected=1. #include <string.h>
* Таймаут = TGC_DIRECT_TIMEOUT_MS. Если цель достигнута → done. #include <limits.h>
*
* Phase 2: последовательный перебор (supernode, затем public/EIM адреса).
* После каждого TIMEOUT — пауза перед следующей попыткой.
*
* Phase 3: последовательный перебор локальных/strict NAT адресов.
* После каждого TIMEOUT — пауза перед следующей попыткой.
*
* После исчерпания Phase 3 → пауза → cycle_restart → Phase 1 (бесконечно).
*
* Цель (tgc_goal_reached): подключены хотя бы к одному суперузлу (node_type==4)
* ЛИБО к 3+ не-мобильным клиентам (device_type != MOBILE, не суперузлам).
* Пока цель не достигнута, успешное подключение НЕ останавливает цикл —
* продолжается перебор следующих кандидатов.
*
* Пауза: активный режим — 1с; фоновый (Android standby) — standby_wait
* (burst → сразу, sleep → до следующего burst).
* При DOWN, если цель перестала выполняться — немедленный cycle_restart
* (уже установленные соединения не трогаются).
* При повторном topo_group_connect_init — destroy + fresh init.
*
* Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями.
* UP: инкремент только для первого соединения к peer_node_id.
* DOWN: декремент только если нет других conn к peer_node_id И нет indirect-путей.
*/
#include "topo_group_connect.h" #include "topo_group_connect.h"
#include "topo_group.h" #include "topo_group.h"
#include "topo_node.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "topo_recovery.h"
#include "../utun_instance.h" #include "../utun_instance.h"
#include "../transport_layer/node_conn_direct.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/mem.h" #include "../lib/mem.h"
#include "../lib/u_async.h" #include "../lib/platform_compat.h"
#include "../lib/ll_queue.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "../chat/chat_event.h" #include "../chat/chat_event.h"
#include "etcp.h"
#ifdef UTUN_HAVE_STANDBY #ifdef UTUN_HAVE_STANDBY
#include "standby.h" #include "standby.h"
#endif #endif
/* ─── внутренние константы ─── */ #define TGC_PARALLEL 128
#define TGC_ID "topo_group_connect" #define TGC_PAUSE_TB 10000
#define TGC_DIRECT_TIMEOUT_MS 2000
#define TGC_PHASE_ONE 0 struct tgc_candidate {
#define TGC_PHASE_TWO 1 uint64_t node_id, started;
#define TGC_PHASE_THREE 2 struct TOPO_PEER_REQUEST* request;
#define TGC_PHASE_DONE 3 uint8_t priority, tried, manual;
#define TGC_PHASE_PAUSE 4 };
#define TGC_MAX_HANDLES 128
#define TGC_PAUSE_ACTIVE_TB 10000 /* 1s пауза в активном режиме */
#define TGC_GOAL_SUPERNODES 1 /* цель: подключены хотя бы к одному суперузлу */
#define TGC_GOAL_NONMOBILE 3 /* цель: подключены к 3+ не-мобильным клиентам (не суперузлам) */
struct TOPO_GROUP_CONNECT { struct TOPO_GROUP_CONNECT {
struct TOPO_GROUP* group; struct TOPO_GROUP* group;
void* phase_timer; struct tgc_candidate* candidates;
void* pause_timer; /* uasync timeout handle (активный режим) */ size_t count;
void* pause_wait; /* standby_wait handle (фоновый режим, Android) */ void *wake, *timer, *pause_timer, *pause_wait;
uint8_t phase; uint8_t automatic, cycle_done, paused, reload_after_pause;
int pending;
int connected_count;
int active_conn_count;
int cursor;
int tried_super;
uint64_t* candidate_ids;
int candidate_count;
struct NODE_CONN_DIRECT* handles[TGC_MAX_HANDLES];
int handle_count;
struct NODE_CONN_DIRECT* cur_handle; /* handle текущей попытки Phase 2/3 (для закрытия на TIMEOUT) */
}; };
static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc); static void tgc_step(void* arg);
static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc);
static void tgc_phase1_timeout(void* arg);
static void tgc_cycle_restart(struct TOPO_GROUP_CONNECT* gc);
static void tgc_start_pause(struct TOPO_GROUP_CONNECT* gc, void (*cb)(void*), const char* label);
static int tgc_goal_reached(struct TOPO_GROUP_CONNECT* gc);
static void tgc_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg);
/* ─── нотификация GUI о списке узлов в процессе подключения ─── */
/**
* Отправляет в GUI событие CHAT_EVT_CONNECTING_NODES — список узлов, к которым
* сейчас идёт попытка подключения. GUI по нему рисует жёлтый кружок у этих узлов.
* Пустой список (count=0) означает «сейчас ни к кому не подключаемся» —
* все жёлтые кружки снимаются.
*/
static void tgc_notify_connecting(struct TOPO_GROUP_CONNECT* gc, const uint64_t* ids, int count) {
size_t cl = strlen(gc->group->channel_id); if (cl > 255) cl = 255;
size_t sz = 1 + cl + 2 + (size_t)count * 8;
uint8_t* buf = u_malloc(sz); if (!buf) return;
uint8_t* p = buf;
*p++ = (uint8_t)cl; memcpy(p, gc->group->channel_id, cl); p += cl;
uint16_t c = (uint16_t)count; memcpy(p, &c, 2); p += 2;
for (int i = 0; i < count; i++) { memcpy(p, &ids[i], 8); p += 8; }
chat_event_post(gc->group->instance, CHAT_EVT_CONNECTING_NODES, buf, (int)sz);
u_free(buf);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: notify_connecting ch=%s count=%d",
TGC_ID, gc->group->channel_id, count);
}
/* ═══════════════════════════════════════════════════════════════════════
* Условие достижения цели авто-подключения
* ══════════════════════════════════════════════════════════════════════ */
/* Тип устройства пира из handshake (peer_device_type на первом UP-линке). */
static uint8_t tgc_peer_device_type(struct ETCP_CONN* conn) {
if (!conn || !conn->links) return CLIENT_TYPE_SERVER;
return conn->links->peer_device_type;
}
/**
* Достигнута ли цель подключения к группе:
* - подключены хотя бы к одному суперузлу (peers node_type==4), ЛИБО
* - подключены к 3+ не-мобильным клиентам (device_type != MOBILE и не суперузел).
*
* Сканирует senders_list (активные BGP-пиры) с дедупом по node_id.
* Состояние не хранится — пересчитывается на лету в редких точках принятия
* решения (таймаут фазы, каждый успех, каждый down).
*/
static int tgc_goal_reached(struct TOPO_GROUP_CONNECT* gc) {
if (!gc || !gc->group || !gc->group->senders_list) return 0;
struct TOPO_GROUP* group = gc->group;
sqlite3* db = group->instance->topo_sqlite_db;
uint64_t seen[TGC_MAX_HANDLES]; int seen_count = 0;
int supernode_count = 0, nonmobile_count = 0;
struct ll_entry* e = group->senders_list->head;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
struct ETCP_CONN* conn = item ? item->conn : NULL;
if (conn) {
int dup = 0;
for (int i = 0; i < seen_count; i++) if (seen[i] == conn->peer_node_id) { dup = 1; break; }
if (!dup && seen_count < TGC_MAX_HANDLES) {
seen[seen_count++] = conn->peer_node_id;
int is_super = db && topo_node_sqlite_get_node_type(db, group->channel_id, conn->peer_node_id) == 4;
int is_mobile = tgc_peer_device_type(conn) == CLIENT_TYPE_MOBILE;
if (is_super) supernode_count++;
else if (!is_mobile) nonmobile_count++;
}
}
e = e->next;
}
int reached = (supernode_count >= TGC_GOAL_SUPERNODES) || (nonmobile_count >= TGC_GOAL_NONMOBILE);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: goal ch=%s super=%d/%d nonmobile=%d/%d → %s",
TGC_ID, group->channel_id, supernode_count, TGC_GOAL_SUPERNODES,
nonmobile_count, TGC_GOAL_NONMOBILE, reached ? "reached" : "not reached");
return reached;
}
/* ═══════════════════════════════════════════════════════════════════════
* Жизненный цикл
* ══════════════════════════════════════════════════════════════════════ */
/**
* Запуск авто-подключения к узлам CHAT-группы.
* Если модуль уже запущен — останавливает его и стартует заново с нуля.
* Начинает с Phase 1: параллельные попытки ко всем пирам, что были connected
* в прошлой сессии (из БД). Если таких нет — сразу переходит к Phase 2.
*/
int topo_group_connect_init(struct TOPO_GROUP* group) {
if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0])
return -1;
if (group->connect) topo_group_connect_destroy(group);
struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc));
if (!gc) return -1;
gc->group = group; gc->tried_super = 0;
group->connect = gc;
uint64_t* ids = NULL; int count = 0; static void tgc_notify(struct TOPO_GROUP_CONNECT* gc) {
sqlite3* db = group->instance->topo_sqlite_db; size_t length = strlen(gc->group->channel_id), count = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: init ch=%s db=%p", TGC_ID, group->channel_id, (void*)db); for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) count++;
if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) { size_t size = 1 + length + 2 + count * 8;
if (ids) { u_free(ids); ids = NULL; } uint8_t* data = u_malloc(size);
gc->phase = TGC_PHASE_TWO; if (!data) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "connecting notification allocation failed"); return; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s no connected peers → Phase 2", TGC_ID, group->channel_id); data[0] = (uint8_t)length; memcpy(data + 1, gc->group->channel_id, length);
tgc_phase2_try_next(gc); uint16_t n = (uint16_t)count; memcpy(data + 1 + length, &n, 2);
return 0; uint8_t* out = data + 3 + length;
for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) {
memcpy(out, &gc->candidates[i].node_id, 8); out += 8;
} }
chat_event_post(gc->group->instance, CHAT_EVT_CONNECTING_NODES, data, (int)size); u_free(data);
gc->phase = TGC_PHASE_ONE;
int launched = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase 1 launching %d connects ch=%s grp=%016llx", TGC_ID, count, group->channel_id, (unsigned long long)group->group_id);
for (int i = 0; i < count && gc->handle_count < TGC_MAX_HANDLES; i++) {
if (ids[i] == group->instance->node_id) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 skip self 0x%016llx", TGC_ID, (unsigned long long)ids[i]); continue; }
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 connect to 0x%016llx (%d/%d)", TGC_ID, (unsigned long long)ids[i], i + 1, count);
node_conn_direct_open(group->instance, ids[i], tgc_callback, gc,
&gc->handles[gc->handle_count++], NULL);
launched++;
}
gc->pending = launched;
tgc_notify_connecting(gc, ids, gc->handle_count);
u_free(ids);
gc->phase_timer = uasync_set_timeout(group->instance->ua,
TGC_DIRECT_TIMEOUT_MS * 10, gc, tgc_phase1_timeout, "tgc_phase1");
return 0;
} }
/**
* Полная остановка авто-подключения: отменяет все таймеры, закрывает все
* открытые соединения и освобождает память. Вызывается при удалении группы
* и при перезапуске (init/restart).
*/
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; }
if (gc->pause_timer) { uasync_cancel_timeout(group->instance->ua, gc->pause_timer); gc->pause_timer = NULL; }
#ifdef UTUN_HAVE_STANDBY
if (gc->pause_wait) { standby_wait_cancel(gc->pause_wait); gc->pause_wait = NULL; }
#endif
for (int i = 0; i < gc->handle_count; i++) node_conn_direct_close(gc->handles[i]);
if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; }
u_free(gc->candidate_ids);
u_free(gc);
}
/**
* Внешний перезапуск авто-подключения (destroy + init). Используется, например,
* при смене сети, когда надо начать поиск заново. Не трогает модуль, если уже
* есть живые соединения.
*/
void topo_group_connect_restart(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s no gc → init", TGC_ID, group->channel_id); topo_group_connect_init(group); return; }
if (gc->active_conn_count > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip active=%d", TGC_ID, group->channel_id, gc->active_conn_count); return; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s phase=%d", TGC_ID, group->channel_id, gc->phase);
topo_group_connect_destroy(group);
topo_group_connect_init(group);
}
/**
* Сколько сейчас живых соединений с пирами группы. Нужно внешней логике
* (например, chat_sync) для решения, запускать ли переподключение.
*/
int topo_group_connect_active_count(struct TOPO_GROUP* group) { int topo_group_connect_active_count(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect; int count = 0;
return gc ? gc->active_conn_count : 0; for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) {
struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (topo_group_peer_ready(group, peer->node_id)) count++;
} }
return count;
/* ═══════════════════════════════════════════════════════════════════════
* ON UP / DOWN
* ══════════════════════════════════════════════════════════════════════ */
/**
* Обработка появления соединения с пиром. Увеличивает счётчик активных и
* помечает connected=1 в БД (чтобы при следующем старте узел попал в Phase 1).
* Пропускает, если пир уже подключён через другое соединение.
*/
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;
if (peer == group->instance->node_id) return;
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) return;
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++; static int tgc_goal(struct TOPO_GROUP_CONNECT* gc) {
sqlite3* db = group->instance->topo_sqlite_db; int super = 0, other = 0;
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1); struct TOPO_GROUP* group = gc->group;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s grp=%016llx active=%d", for (struct ll_entry* e = group->senders_list->head; e; e = e->next) {
TGC_ID, (unsigned long long)peer, group->channel_id, (unsigned long long)group->group_id, gc->active_conn_count); struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data;
} if (!topo_group_peer_ready(group, peer->node_id)) continue;
if (topo_node_sqlite_get_node_type(group->instance->topo_sqlite_db, group->channel_id, peer->node_id) == 4) super++;
/** else if (!peer->conn->links || peer->conn->links->peer_device_type != CLIENT_TYPE_MOBILE) other++;
* Обработка обрыва соединения. Уменьшает счётчик активных и помечает
* connected=0 в БД. Если живых соединений не осталось — немедленно
* перезапускает цикл поиска с Phase 1.
*/
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;
if (peer == group->instance->node_id) return;
struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) return;
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; return super >= 1 || other >= 3;
} }
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, peer); static struct tgc_candidate* tgc_add(struct TOPO_GROUP_CONNECT* gc, uint64_t id, uint8_t priority, int manual) {
if (nq && nq->paths) { for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].node_id == id) return &gc->candidates[i];
struct ll_entry* pe = nq->paths->head; struct tgc_candidate* list = u_realloc(gc->candidates, (gc->count + 1) * sizeof(*list));
while (pe) { if (!list) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group candidate allocation failed"); return NULL; }
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe; gc->candidates = list;
if (path->conn && path->conn != conn) { list[gc->count] = (struct tgc_candidate){ .node_id = id, .priority = priority, .manual = manual };
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx has indirect path, skip", return &list[gc->count++];
TGC_ID, (unsigned long long)peer);
return;
}
pe = pe->next;
}
} }
gc->active_conn_count--; static void tgc_load(struct TOPO_GROUP_CONNECT* gc) {
sqlite3* db = group->instance->topo_sqlite_db; size_t retained = 0;
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0); for (size_t i = 0; i < gc->count; i++) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d", struct tgc_candidate* c = &gc->candidates[i];
TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); if (c->manual && c->request) gc->candidates[retained++] = *c;
if (!tgc_goal_reached(gc)) { else topo_group_peer_close(c->request);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN goal not reached — cycle restart ch=%s", TGC_ID, group->channel_id);
tgc_cycle_restart(gc);
} }
gc->count = retained; gc->cycle_done = 0;
sqlite3* db = gc->group->instance->topo_sqlite_db;
for (int priority = 0; priority < 4; priority++) {
uint64_t* ids = NULL; int count = 0, rc;
if (priority == 0) rc = topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count);
else if (priority == 1) rc = topo_node_sqlite_get_supernode_peers(db, gc->group->channel_id, &ids, &count);
else if (priority == 2) rc = topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &ids, &count, gc->group->instance->node_id);
else rc = topo_node_sqlite_get_local_peers(db, gc->group->channel_id, &ids, &count, gc->group->instance->node_id);
if (rc < 0) DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot load group candidates group=%016llx priority=%d",
(unsigned long long)gc->group->group_id, priority);
else for (int i = 0; i < count; i++)
if (ids[i] != gc->group->instance->node_id && !tgc_add(gc, ids[i], priority, 0)) break;
u_free(ids);
} }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect cycle: group=%016llx candidates=%zu",
/* ═══════════════════════════════════════════════════════════════════════ (unsigned long long)gc->group->group_id, gc->count);
* Единый коллбэк для всех фаз
* ══════════════════════════════════════════════════════════════════════ */
/**
* Единый коллбэк на результат каждой попытки подключения (успех/таймаут/обрыв).
* В Phase 1 — просто считает результаты (итог подводит tgc_phase1_timeout).
* В Phase 2/3 — успех проверяет цель: достигнута → done, иначе следующий кандидат;
* провал закрывает handle (чтобы не оставить залипшее состояние) и ставит паузу.
*/
/* Передача владения: после topo_group_new_conn (который взял свой NCD-handle)
* закрываем connect-фазу handle, чтобы владельцем conn остался один BGP-хэндл
* (в senders_list), без двойного счёта. */
static void tgc_release_handle(struct TOPO_GROUP_CONNECT* gc, struct NODE_CONN_DIRECT* h) {
if (!gc || !h) return;
for (int i = 0; i < gc->handle_count; i++) {
if (gc->handles[i] == h) { gc->handles[i] = NULL; node_conn_direct_close(h); return; }
}
if (gc->cur_handle == h) { gc->cur_handle = NULL; node_conn_direct_close(h); return; }
node_conn_direct_close(h);
} }
static void tgc_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { static void tgc_cancel_pause(struct TOPO_GROUP_CONNECT* gc) {
struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; if (gc->pause_timer) uasync_cancel_timeout(gc->group->instance->ua, gc->pause_timer);
struct ETCP_CONN* conn = node_conn_direct_get_conn(h); #ifdef UTUN_HAVE_STANDBY
uint64_t node_id = node_conn_direct_node_id(h); if (gc->pause_wait) standby_wait_cancel(gc->pause_wait);
#endif
int ok = (event == NCD_EVENT_UP); gc->pause_timer = NULL; gc->pause_wait = NULL; gc->paused = 0;
if (ok) {
/* запуск BGP: узел соединён с группой (topo_group_new_conn берёт владение NCD-handle) */
if (conn) topo_group_new_conn(gc->group, conn);
tgc_release_handle(gc, h);
} else {
// Каждая попытка даёт один результат. Поздние DOWN/CLOSED не должны уменьшать pending повторно.
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(gc->group, node_id);
if (nq) nq->handle = NULL;
tgc_release_handle(gc, h);
} }
switch (gc->phase) { static void tgc_resume(void* arg) {
case TGC_PHASE_ONE: struct TOPO_GROUP_CONNECT* gc = arg;
gc->pending--; gc->pause_timer = NULL; gc->pause_wait = NULL; gc->paused = 0;
if (ok) gc->connected_count++; if (gc->reload_after_pause) tgc_load(gc);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 result node=0x%016llx %s pending=%d connected=%d", topo_group_connect_changed(gc->group);
TGC_ID, (unsigned long long)node_id, ok ? "OK" : "FAIL", gc->pending, gc->connected_count);
break;
case TGC_PHASE_TWO:
if (ok) {
if (tgc_goal_reached(gc)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — goal reached, done", TGC_ID, (unsigned long long)node_id);
gc->phase = TGC_PHASE_DONE;
tgc_notify_connecting(gc, NULL, 0);
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — goal not reached, try next", TGC_ID, (unsigned long long)node_id);
tgc_phase2_try_next(gc);
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — pause then try next", TGC_ID, (unsigned long long)node_id);
if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; }
tgc_start_pause(gc, (void(*)(void*))tgc_phase2_try_next, "tgc_p2_try");
}
break;
case TGC_PHASE_THREE:
if (ok) {
if (tgc_goal_reached(gc)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 connected to 0x%016llx — goal reached, done", TGC_ID, (unsigned long long)node_id);
gc->phase = TGC_PHASE_DONE;
tgc_notify_connecting(gc, NULL, 0);
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 connected to 0x%016llx — goal not reached, try next", TGC_ID, (unsigned long long)node_id);
tgc_phase3_try_next(gc);
}
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 0x%016llx FAIL — pause then try next", TGC_ID, (unsigned long long)node_id);
if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; }
tgc_start_pause(gc, (void(*)(void*))tgc_phase3_try_next, "tgc_p3_try");
}
break;
}
} }
/* ═══════════════════════════════════════════════════════════════════════ static void tgc_pause(struct TOPO_GROUP_CONNECT* gc, int reload) {
* Пауза: 1s ACTIVE / 30s STANDBY gc->paused = 1; gc->reload_after_pause = reload;
* ══════════════════════════════════════════════════════════════════════ */
/**
* Пауза между попытками подключения. Длительность зависит от активности
* приложения: 1 секунда когда приложение активно, 30 секунд когда в фоне
* (экран выключен). После паузы вызывает cb (следующая попытка или перезапуск
* цикла). На время паузы уведомляет GUI пустым списком (жёлтые кружки гаснут).
*/
static void tgc_start_pause(struct TOPO_GROUP_CONNECT* gc, void (*cb)(void*), const char* label) {
gc->phase = TGC_PHASE_PAUSE;
tgc_notify_connecting(gc, NULL, 0);
#ifdef UTUN_HAVE_STANDBY #ifdef UTUN_HAVE_STANDBY
if (standby_is_enabled()) { if (standby_is_enabled()) {
/* фоновый режим: burst → standby_wait будит сразу, sleep → до следующего burst */ gc->pause_wait = standby_wait(gc, tgc_resume);
gc->pause_wait = standby_wait(gc, cb); if (gc->pause_wait) return;
gc->pause_timer = NULL; DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot wait for standby; using group retry timer");
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: pause (standby) for ch=%s", TGC_ID, gc->group->channel_id);
return;
} }
#endif #endif
gc->pause_timer = uasync_set_timeout(gc->group->instance->ua, TGC_PAUSE_ACTIVE_TB, gc, cb, label); gc->pause_timer = uasync_set_timeout(gc->group->instance->ua, TGC_PAUSE_TB, gc, tgc_resume, "group_connect_pause");
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: pause %dms (active) for ch=%s", if (!gc->pause_timer) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group retry timer allocation failed"); gc->paused = 0; }
TGC_ID, TGC_PAUSE_ACTIVE_TB / 10, gc->group->channel_id); }
static int tgc_open(struct TOPO_GROUP_CONNECT* gc, struct tgc_candidate* c) {
c->tried = 1; c->started = get_time_tb();
int result = topo_group_peer_open(gc->group, c->node_id, &c->request);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect attempt: group=%016llx node=%016llx priority=%u manual=%u result=%d",
(unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, c->priority, c->manual, result);
if (result == 0) topo_group_connect_changed(gc->group);
return result;
}
static uint64_t tgc_deadline(struct tgc_candidate* c, enum topo_peer_phase phase) {
return phase == TOPO_PEER_CONNECTING ? c->started + TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10ULL :
topo_group_peer_progress(c->request) + TOPO_RECOVERY_SYNC_TIMEOUT_MS * 10ULL;
}
static void tgc_timeout(void* arg) {
struct TOPO_GROUP_CONNECT* gc = arg; gc->timer = NULL; tgc_step(gc);
}
static void tgc_step(void* arg) {
struct TOPO_GROUP_CONNECT* gc = arg;
if (gc->wake) { uasync_call_soon_cancel(gc->group->instance->ua, gc->wake); gc->wake = NULL; }
if (gc->timer) { uasync_cancel_timeout(gc->group->instance->ua, gc->timer); gc->timer = NULL; }
uint64_t now = get_time_tb();
int changed = 0, active = 0, failed = 0;
for (size_t i = 0; i < gc->count; i++) {
struct tgc_candidate* c = &gc->candidates[i];
if (!c->request) continue;
enum topo_peer_phase phase = topo_group_peer_phase(c->request);
if (phase == TOPO_PEER_READY || phase == TOPO_PEER_FAILED || now >= tgc_deadline(c, phase)) {
if (phase != TOPO_PEER_READY) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "group connect failed: group=%016llx node=%016llx phase=%d",
(unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, phase);
if (!c->manual && c->priority) failed = 1;
} else DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect READY: group=%016llx node=%016llx",
(unsigned long long)gc->group->group_id, (unsigned long long)c->node_id);
struct TOPO_PEER_REQUEST* request = c->request; c->request = NULL;
topo_group_peer_close(request); changed = 1;
} else if (!c->manual) active++;
}
if (gc->automatic && tgc_goal(gc)) {
tgc_cancel_pause(gc);
if (!gc->cycle_done) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect goal reached: group=%016llx", (unsigned long long)gc->group->group_id);
for (size_t i = 0; i < gc->count; i++) if (!gc->candidates[i].manual && gc->candidates[i].request) {
struct TOPO_PEER_REQUEST* request = gc->candidates[i].request; gc->candidates[i].request = NULL;
topo_group_peer_close(request); changed = 1;
}
gc->cycle_done = 1;
}
} else if (gc->automatic && !gc->paused) {
if (gc->cycle_done) { tgc_load(gc); active = 0; }
if (failed) tgc_pause(gc, 0);
else {
/* Priority 0 — parallel historical peers; all other attempts are sequential. */
int available = 0;
for (size_t i = 0; i < gc->count; i++) {
struct tgc_candidate* c = &gc->candidates[i];
if (c->manual || c->tried) continue;
available = 1;
if (active && (c->priority || active >= TGC_PARALLEL)) break;
if (tgc_open(gc, c) == 0) { active++; changed = 1; }
}
if (!available && !active) tgc_pause(gc, 1);
}
}
uint64_t earliest = UINT64_MAX;
for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) {
uint64_t deadline = tgc_deadline(&gc->candidates[i], topo_group_peer_phase(gc->candidates[i].request));
if (deadline < earliest) earliest = deadline;
}
if (earliest != UINT64_MAX) {
uint64_t delay = earliest > now ? earliest - now : 1;
gc->timer = uasync_set_timeout(gc->group->instance->ua, delay > INT_MAX ? INT_MAX : (int)delay, gc, tgc_timeout, "group_connect");
if (!gc->timer) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect timer allocation failed");
}
if (changed) tgc_notify(gc);
}
static void tgc_wake(void* arg) {
struct TOPO_GROUP_CONNECT* gc = arg; gc->wake = NULL; tgc_step(gc);
}
void topo_group_connect_changed(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group ? group->connect : NULL;
if (!gc || group->stopping || gc->wake) return;
gc->wake = uasync_call_soon(group->instance->ua, gc, tgc_wake);
if (!gc->wake) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect wake allocation failed");
} }
/* ═══════════════════════════════════════════════════════════════════════ void topo_group_connect_destroy(struct TOPO_GROUP* group) {
* Phase 1 timeout struct TOPO_GROUP_CONNECT* gc = group ? group->connect : NULL;
* ══════════════════════════════════════════════════════════════════════ */ if (!gc) return;
group->connect = NULL;
/** if (gc->wake) uasync_call_soon_cancel(group->instance->ua, gc->wake);
* Таймаут параллельных попыток Phase 1 (2 секунды). Закрывает неудачные if (gc->timer) uasync_cancel_timeout(group->instance->ua, gc->timer);
* handles и сбрасывает connected в БД. Если цель достигнута — цикл завершён; tgc_cancel_pause(gc);
* иначе переходим к Phase 2 (дополнительные попытки). for (size_t i = 0; i < gc->count; i++) topo_group_peer_close(gc->candidates[i].request);
*/ u_free(gc->candidates); u_free(gc);
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;
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++) {
int connected = 0;
for (int j = 0; j < gc->handle_count; j++) {
struct NODE_CONN_DIRECT* h = gc->handles[j];
if (h && node_conn_direct_get_conn(h) != NULL) { connected = 1; break; }
}
if (!connected) topo_node_sqlite_set_connected(db, gc->group->channel_id, ids[i], 0);
}
u_free(ids);
int remaining = 0;
for (int i = 0; i < gc->handle_count; i++) {
struct NODE_CONN_DIRECT* h = gc->handles[i];
if (h && node_conn_direct_get_conn(h) != NULL) {
gc->handles[remaining++] = h;
} else {
if (h) node_conn_direct_close(h);
}
}
gc->handle_count = remaining;
if (tgc_goal_reached(gc)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — goal reached (%d connected), skipping Phase 2", TGC_ID, gc->connected_count);
gc->phase = TGC_PHASE_DONE;
tgc_notify_connecting(gc, NULL, 0);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — goal not reached (%d connected) → Phase 2", TGC_ID, gc->connected_count);
gc->phase = TGC_PHASE_TWO;
tgc_phase2_try_next(gc);
} }
/* ═══════════════════════════════════════════════════════════════════════ int topo_group_connect_init(struct TOPO_GROUP* group) {
* Phase 2 if (!group || group->stopping || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0]) {
* ══════════════════════════════════════════════════════════════════════ */ DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot start group connect without a CHAT group"); return -1;
/**
* Последовательный перебор узлов с публичными/EIM адресами (сначала supernode,
* затем обычные публичные). Пробует следующий узел и ставит ему жёлтый кружок.
* Когда все перебраны — переходит к Phase 3.
*/
static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) {
gc->pause_timer = NULL; gc->pause_wait = NULL;
if (gc->phase == TGC_PHASE_PAUSE) gc->phase = TGC_PHASE_TWO;
if (gc->phase != TGC_PHASE_TWO) return;
while (1) {
if (gc->candidate_count == 0) {
sqlite3* db = gc->group->instance->topo_sqlite_db;
if (gc->tried_super == 0) {
topo_node_sqlite_get_supernode_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count);
gc->tried_super = 1; gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 supernode round — loaded %d for ch=%s",
TGC_ID, gc->candidate_count, gc->group->channel_id);
} else if (gc->tried_super == 1) {
topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count,
gc->group->instance->node_id);
gc->tried_super = 2; gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 public round — loaded %d for ch=%s",
TGC_ID, gc->candidate_count, gc->group->channel_id);
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id);
gc->phase = TGC_PHASE_THREE; gc->candidate_count = 0;
tgc_phase3_try_next(gc);
return;
}
}
if (gc->cursor < gc->candidate_count) {
uint64_t nid = gc->candidate_ids[gc->cursor++];
if (nid == gc->group->instance->node_id) continue;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d) super_round=%d", TGC_ID,
(unsigned long long)nid, gc->cursor, gc->candidate_count, gc->tried_super);
if (node_conn_direct_open(gc->group->instance, nid, tgc_callback, gc, &gc->cur_handle, NULL) < 0) continue;
tgc_notify_connecting(gc, &nid, 1);
return;
}
if (gc->candidate_ids) { u_free(gc->candidate_ids); gc->candidate_ids = NULL; }
gc->candidate_count = 0;
}
} }
topo_group_connect_destroy(group);
/* ═══════════════════════════════════════════════════════════════════════ struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc));
* Phase 3 — локальные/strict NAT адреса if (!gc) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect allocation failed"); return -1; }
* ══════════════════════════════════════════════════════════════════════ */ gc->group = group; gc->automatic = 1; group->connect = gc;
tgc_load(gc); topo_group_connect_changed(group); return 0;
/**
* Последовательный перебор узлов с локальными/strict NAT адресами (как Phase 2,
* но для узлов, доступных только в локальной сети). Когда все перебраны —
* пауза и полный перезапуск цикла с Phase 1 (бесконечно, пока не подключимся).
*/
static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc) {
gc->pause_timer = NULL; gc->pause_wait = NULL;
if (gc->phase == TGC_PHASE_PAUSE) gc->phase = TGC_PHASE_THREE;
if (gc->phase != TGC_PHASE_THREE) return;
if (gc->candidate_count == 0) {
sqlite3* db = gc->group->instance->topo_sqlite_db;
topo_node_sqlite_get_local_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count,
gc->group->instance->node_id);
gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 loaded %d local 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++];
if (nid == gc->group->instance->node_id) continue;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 trying 0x%016llx (%d/%d)", TGC_ID,
(unsigned long long)nid, gc->cursor, gc->candidate_count);
if (node_conn_direct_open(gc->group->instance, nid, tgc_callback, gc, &gc->cur_handle, NULL) < 0) continue;
tgc_notify_connecting(gc, &nid, 1);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 exhausted — pause then cycle restart, ch=%s", TGC_ID, gc->group->channel_id);
tgc_start_pause(gc, (void(*)(void*))tgc_cycle_restart, "tgc_p3_restart");
} }
/* ═══════════════════════════════════════════════════════════════════════ void topo_group_connect_restart(struct TOPO_GROUP* group) {
* Полный перезапуск цикла с Phase 1 (без destroy/init) if (!group) return;
* ══════════════════════════════════════════════════════════════════════ */ if (topo_group_connect_active_count(group)) return;
topo_group_connect_init(group);
/**
* Полный перезапуск цикла с Phase 1 без уничтожения модуля. Закрывает все
* открытые соединения, сбрасывает состояние и заново запускает параллельные
* попытки ко всем connected-пирам из БД. Вызывается при обрыве всех соединений
* и после исчерпания Phase 3.
*/
static void tgc_cycle_restart(struct TOPO_GROUP_CONNECT* gc) {
gc->pause_timer = NULL; gc->pause_wait = NULL;
if (gc->phase_timer) { uasync_cancel_timeout(gc->group->instance->ua, gc->phase_timer); gc->phase_timer = NULL; }
for (int i = 0; i < gc->handle_count; i++) node_conn_direct_close(gc->handles[i]);
if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; }
gc->handle_count = 0; gc->pending = 0; gc->connected_count = 0;
gc->cursor = 0; gc->tried_super = 0;
u_free(gc->candidate_ids); gc->candidate_ids = NULL; gc->candidate_count = 0;
uint64_t* ids = NULL; int count = 0;
sqlite3* db = gc->group->instance->topo_sqlite_db;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: cycle_restart ch=%s active=%d",
TGC_ID, gc->group->channel_id, gc->active_conn_count);
if (!db || topo_node_sqlite_get_connected_peers(db, gc->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: cycle_restart ch=%s no connected peers → Phase 2", TGC_ID, gc->group->channel_id);
tgc_phase2_try_next(gc);
return;
} }
gc->phase = TGC_PHASE_ONE; void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
int launched = 0; if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !conn) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: cycle_restart Phase 1 launching %d connects ch=%s", if (group->instance->topo_sqlite_db && topo_group_peer_ready(group, conn->peer_node_id))
TGC_ID, count, gc->group->channel_id); topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 1);
for (int i = 0; i < count && gc->handle_count < TGC_MAX_HANDLES; i++) { topo_group_connect_changed(group);
if (ids[i] == gc->group->instance->node_id) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 skip self 0x%016llx", TGC_ID, (unsigned long long)ids[i]); continue; }
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 connect to 0x%016llx (%d/%d)", TGC_ID, (unsigned long long)ids[i], i + 1, count);
node_conn_direct_open(gc->group->instance, ids[i], tgc_callback, gc,
&gc->handles[gc->handle_count++], NULL);
launched++;
}
gc->pending = launched;
tgc_notify_connecting(gc, ids, gc->handle_count);
u_free(ids);
gc->phase_timer = uasync_set_timeout(gc->group->instance->ua,
TGC_DIRECT_TIMEOUT_MS * 10, gc, tgc_phase1_timeout, "tgc_phase1");
} }
/* ═══════════════════════════════════════════════════════════════════════ void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
* Однократное подключение к конкретному узлу (ручной "Connect" из GUI) if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !conn) return;
* ══════════════════════════════════════════════════════════════════════ */ if (group->instance->topo_sqlite_db)
topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 0);
struct tgc_once_ctx { topo_group_connect_changed(group);
struct TOPO_GROUP* group;
uint64_t node_id;
struct NODE_CONN_DIRECT* h;
};
static void tgc_once_cb(struct NODE_CONN_DIRECT* h, enum ncd_event ev, void* arg) {
struct tgc_once_ctx* c = (struct tgc_once_ctx*)arg;
struct ETCP_CONN* conn = node_conn_direct_get_conn(h);
if (ev == NCD_EVENT_UP && conn) {
node_conn_direct_set_callback(h, NULL, NULL);
topo_group_new_conn(c->group, conn);
node_conn_direct_close(h); /* владение перешло к topo_group_new_conn */
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect UP node=0x%016llx ch=%s",
TGC_ID, (unsigned long long)c->node_id, c->group->channel_id);
} else if (ev == NCD_EVENT_TIMEOUT) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "%s: once-connect TIMEOUT node=0x%016llx ch=%s",
TGC_ID, (unsigned long long)c->node_id, c->group->channel_id);
node_conn_direct_close(h);
} else if (ev == NCD_EVENT_DOWN || ev == NCD_EVENT_CLOSED) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "%s: once-connect DOWN node=0x%016llx ch=%s",
TGC_ID, (unsigned long long)c->node_id, c->group->channel_id);
node_conn_direct_close(h);
}
u_free(c);
} }
int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id) { int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id) {
if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !node_id) return -1; if (!group || group->stopping || group->group_type != TOPO_GROUP_TYPE_CHAT || !node_id || node_id == group->instance->node_id) {
if (node_id == group->instance->node_id) return -1; DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid manual group connect"); return -1;
struct ETCP_CONN* ex = instance_find_conn(group->instance, node_id);
if (ex && ex->links_up && ex->initialized) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect node=0x%016llx already connected — add to group ch=%s",
TGC_ID, (unsigned long long)node_id, group->channel_id);
topo_group_new_conn(group, ex);
return 0;
} }
if (!group->connect) {
struct tgc_once_ctx* c = u_calloc(1, sizeof(*c)); group->connect = u_calloc(1, sizeof(*group->connect));
if (!group->connect) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "manual group connect allocation failed"); return -1; }
group->connect->group = group;
}
struct TOPO_GROUP_CONNECT* gc = group->connect;
struct tgc_candidate* c = tgc_add(gc, node_id, 0, 1);
if (!c) return -1; if (!c) return -1;
c->group = group; c->node_id = node_id; c->manual = 1;
if (c->request) return 0;
int r = node_conn_direct_open(group->instance, node_id, tgc_once_cb, c, &c->h, NULL); int result = tgc_open(gc, c); tgc_notify(gc); return result;
if (r == NCD_ERR || !c->h) { u_free(c); return -1; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect node=0x%016llx ch=%s started (rc=%d)",
TGC_ID, (unsigned long long)node_id, group->channel_id, r);
return 0;
} }

36
src/routing_layer/topo_group_connect.h

@ -1,26 +1,11 @@
/** /* CHAT connection policy over cancellable group requests.
* @file topo_group_connect.h * Phase 1: historical connected peers in parallel (up to 128).
* @brief Авто-подключение к узлам группы — бесконечный цикл до цели. * Phase 2: sequential supernode/public peers. Phase 3: local peers.
* * Each node is tried once per cycle. CONNECTING: 2s; SYNCING: 5s without progress.
* При старте CHAT-группы (topo_group_connect_init): * Goal: one READY supernode or three READY nonmobile peers. After exhaustion,
* Phase 1 — параллельный запуск ко всем пирам с connected=1 (таймаут 2с). * retry after 1s (or standby burst). Existing READY group sessions survive
* Если цель достигнута → done. * planner restart/stop. Manual attempts are tracked and cancelled with the group.
* Phase 2 — последовательный перебор (supernode, затем public/EIM). * active_count is derived from the current group sessions, not a separate counter.
* После каждого TIMEOUT — пауза перед следующей попыткой.
* Phase 3 — последовательный перебор локальных/strict NAT.
* После каждого TIMEOUT — пауза перед следующей попыткой.
* После исчерпания Phase 3 → пауза → cycle_restart → Phase 1 (бесконечно).
*
* Цель (tgc_goal_reached): подключены хотя бы к одному суперузлу (node_type==4)
* ЛИБО к 3+ не-мобильным клиентам (device_type != MOBILE, не суперузлам).
* Пока цель не достигнута, успех НЕ останавливает цикл.
*
* Пауза: активный режим — 1с; фоновый (Android standby) — standby_wait
* (burst → сразу, sleep → до следующего burst).
* При DOWN, если цель перестала выполняться — немедленный cycle_restart.
* При повторном вызове topo_group_connect_init — destroy + fresh init.
*
* active_conn_count — количество уникальных peer-узлов с живыми соединениями.
*/ */
#ifndef TOPO_GROUP_CONNECT_H #ifndef TOPO_GROUP_CONNECT_H
#define TOPO_GROUP_CONNECT_H #define TOPO_GROUP_CONNECT_H
@ -32,14 +17,13 @@ struct ETCP_CONN;
int topo_group_connect_init(struct TOPO_GROUP* group); int topo_group_connect_init(struct TOPO_GROUP* group);
void topo_group_connect_destroy(struct TOPO_GROUP* group); void topo_group_connect_destroy(struct TOPO_GROUP* group);
void topo_group_connect_changed(struct TOPO_GROUP* group);
void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn); 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); void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
int topo_group_connect_active_count(struct TOPO_GROUP* group); int topo_group_connect_active_count(struct TOPO_GROUP* group);
void topo_group_connect_restart(struct TOPO_GROUP* group); void topo_group_connect_restart(struct TOPO_GROUP* group);
/* Однократная попытка подключения к конкретному узлу группы (без бесконечного /* One manual request; shares the group session, no NCD ownership handoff on UP. */
* цикла). При успехе соединение добавляется в BGP-группу (topo_group_new_conn),
* handle сохраняется в nq->handle. Используется для ручного "Connect" из GUI. */
int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id); int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id);
#endif #endif

49
src/routing_layer/topo_group_invite.c

@ -1,3 +1,4 @@
#include "topo_group_connect.h"
/** /**
* @file topo_group_invite.c * @file topo_group_invite.c
* @brief Invite/join к каналу через прямое (ncd) подключение + проверка членства. * @brief Invite/join к каналу через прямое (ncd) подключение + проверка членства.
@ -187,7 +188,7 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz
{ {
struct TOPO_GROUP* bgp_grp = topo_groups_find(inst->topo_groups, req->group_id); struct TOPO_GROUP* bgp_grp = topo_groups_find(inst->topo_groups, req->group_id);
if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT) if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT)
topo_group_new_conn(bgp_grp, conn); topo_group_connect_node_once(bgp_grp, conn->peer_node_id);
} }
/* ── сериализуем ответ ── */ /* ── сериализуем ответ ── */
@ -342,11 +343,10 @@ static void tgi_handle_info_resp(struct ETCP_CONN* conn, const uint8_t* data, si
{ {
struct TOPO_GROUP* bgp_grp = topo_groups_find(inv->inst->topo_groups, inv->group_id); struct TOPO_GROUP* bgp_grp = topo_groups_find(inv->inst->topo_groups, inv->group_id);
if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT) if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT)
topo_group_new_conn(bgp_grp, conn); topo_group_connect_node_once(bgp_grp, conn->peer_node_id);
} }
/* ── успех: владение conn перешло к topo_group_new_conn (item->handle), /* Bootstrap приглашения завершён; запрос групповой сессии уже держит свой интерес. */
закрываем connect-фазу handle и доставляем JOIN ── */
{ {
uint64_t node_id = inv->node_id, group_id = inv->group_id; uint64_t node_id = inv->node_id, group_id = inv->group_id;
tgi_cb_t cb = inv->cb; void* cb_arg = inv->cb_arg; tgi_cb_t cb = inv->cb; void* cb_arg = inv->cb_arg;
@ -509,23 +509,6 @@ int topo_group_invite_join(struct UTUN_INSTANCE* inst, uint64_t group_id,
/* ═══════════ инвайтера сторона (direct, без INVITE_INFO) ═══════════ */ /* ═══════════ инвайтера сторона (direct, без INVITE_INFO) ═══════════ */
struct tgi_conn_ctx { struct UTUN_INSTANCE* inst; uint64_t node_id, group_id; };
static void tgi_to_channel_cb(struct NODE_CONN_DIRECT* h, enum ncd_event ev, void* arg) {
struct tgi_conn_ctx* c = (struct tgi_conn_ctx*)arg;
if (ev == NCD_EVENT_UP) {
struct ETCP_CONN* conn = node_conn_direct_get_conn(h);
struct TOPO_GROUP* g = topo_groups_find(c->inst->topo_groups, c->group_id);
/* владение переходит к topo_group_new_conn (item->handle), connect-фазу закрываем */
if (g && conn) topo_group_new_conn(g, conn);
node_conn_direct_close(h);
u_free(c);
} else if (ev == NCD_EVENT_TIMEOUT || ev == NCD_EVENT_CLOSED) {
node_conn_direct_force_close(h);
u_free(c);
}
}
int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id, int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id,
struct TOPO_NODE* ni, uint64_t node_id) { struct TOPO_NODE* ni, uint64_t node_id) {
if (!inst || !inst->topo_groups) return -1; if (!inst || !inst->topo_groups) return -1;
@ -537,23 +520,11 @@ int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id,
} }
if (!node_id || node_id == inst->node_id) return -1; if (!node_id || node_id == inst->node_id) return -1;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id); struct NODE_CONN_DIRECT* seed = NULL;
if (!nq) { if (ni && node_conn_direct_open_node(inst, node_id, NULL, NULL, &seed, ni, NULL) == NCD_ERR) {
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_GROUP_NODE)); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot seed invited peer addresses node=%016llx", (unsigned long long)node_id); return -1;
if (!qe) return -1;
nq = (struct TOPO_GROUP_NODE*)qe;
memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry));
nq->node_id = node_id;
queue_data_put_with_index(group->nodes, &nq->ll);
} }
int result = topo_group_connect_node_once(group, node_id);
struct tgi_conn_ctx* c = u_calloc(1, sizeof(*c)); node_conn_direct_close(seed);
if (!c) return -1; return result;
c->inst = inst; c->node_id = node_id; c->group_id = group_id;
struct NODE_CONN_DIRECT* h = NULL;
int r = ni ? node_conn_direct_open_node(inst, node_id, tgi_to_channel_cb, c, &h, ni, NULL)
: node_conn_direct_open(inst, node_id, tgi_to_channel_cb, c, &h, NULL);
if (r == NCD_ERR || !h) { u_free(c); return -1; }
return 0;
} }

39
src/transport_layer/etcp_connections.c

@ -2770,20 +2770,8 @@ int init_sockets(struct UTUN_INSTANCE* instance) {
return 0; // All OK return 0; // All OK
} }
/* Конфигурация владеет только NCD handles; транспортные линки создаёт NCD. */ /* Адреса конфигурации добавляются в NCD при подготовке. Постоянным владельцем
static void config_client_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) { * транспорта сразу становится группа; адаптер держит только запрос участия. */
struct UTUN_INSTANCE* instance = arg;
if (event == NCD_EVENT_UP) {
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
if (group && topo_group_join_peer(group, handle) < 0)
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "configured peer join failed node=%016llx",
(unsigned long long)node_conn_direct_node_id(handle));
} else if (event == NCD_EVENT_TIMEOUT || event == NCD_EVENT_CLOSED) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "configured peer node=%016llx event=%d",
(unsigned long long)node_conn_direct_node_id(handle), event);
}
}
static struct NODE_CONN_DIRECT* config_client_open(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client) { static struct NODE_CONN_DIRECT* config_client_open(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client) {
struct TOPO_NODE ni = { .node_name = client->name }; struct TOPO_NODE ni = { .node_name = client->name };
if (sc_hex_to_binary(client->peer_public_key_hex, ni.public_key, SC_PUBKEY_SIZE) != SC_OK) { if (sc_hex_to_binary(client->peer_public_key_hex, ni.public_key, SC_PUBKEY_SIZE) != SC_OK) {
@ -2839,41 +2827,42 @@ fail:
static void config_handles_close(struct CONFIG_CONN_HANDLE* head) { static void config_handles_close(struct CONFIG_CONN_HANDLE* head) {
while (head) { while (head) {
struct CONFIG_CONN_HANDLE* next = head->next; struct CONFIG_CONN_HANDLE* next = head->next;
node_conn_direct_close(head->handle); u_free(head); head = next; topo_group_peer_close(head->request); u_free(head); head = next;
} }
} }
int init_client_connections(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* clients) { int init_client_connections(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* clients) {
if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "configure clients without instance"); return -1; } if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "configure clients without instance"); return -1; }
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
if (clients && !group) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "UTUN group is not running"); return -1; }
struct CONFIG_CONN_HANDLE* updated = NULL; struct CONFIG_CONN_HANDLE* updated = NULL;
int count = 0; int count = 0;
for (struct CFG_CLIENT* client = clients; client; client = client->next) { for (struct CFG_CLIENT* client = clients; client; client = client->next) {
struct CONFIG_CONN_HANDLE* owner = u_calloc(1, sizeof(*owner)); struct CONFIG_CONN_HANDLE* owner = u_calloc(1, sizeof(*owner));
if (!owner) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client handle allocation failed"); goto fail; } if (!owner) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client handle allocation failed"); goto fail; }
owner->handle = config_client_open(instance, client); struct NODE_CONN_DIRECT* seed = config_client_open(instance, client);
if (!owner->handle) { u_free(owner); goto fail; } if (!seed) { u_free(owner); goto fail; }
owner->node_id = node_conn_direct_node_id(owner->handle); owner->node_id = node_conn_direct_node_id(seed);
int opened = topo_group_peer_open(group, owner->node_id, &owner->request);
node_conn_direct_close(seed);
if (opened < 0) { u_free(owner); goto fail; }
snprintf(owner->name, sizeof(owner->name), "%s", client->name); snprintf(owner->name, sizeof(owner->name), "%s", client->name);
owner->next = updated; updated = owner; count++; owner->next = updated; updated = owner; count++;
} }
struct CONFIG_CONN_HANDLE* old = instance->config_conn_handles; struct CONFIG_CONN_HANDLE* old = instance->config_conn_handles;
instance->config_conn_handles = updated; instance->config_conn_handles = updated;
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
for (struct CONFIG_CONN_HANDLE* owner = old; owner; owner = owner->next) { for (struct CONFIG_CONN_HANDLE* owner = old; owner; owner = owner->next) {
int retained = 0; int retained = 0;
for (struct CONFIG_CONN_HANDLE* next = updated; next; next = next->next) for (struct CONFIG_CONN_HANDLE* next = updated; next; next = next->next)
if (next->node_id == owner->node_id) { retained = 1; break; } if (next->node_id == owner->node_id) { retained = 1; break; }
if (!retained && group && node_conn_direct_get_conn(owner->handle)) topo_group_leave_peer(group, owner->handle); if (!retained && group) topo_group_peer_leave(group, owner->node_id);
}
for (struct CONFIG_CONN_HANDLE* owner = updated; owner; owner = owner->next) {
node_conn_direct_set_callback(owner->handle, config_client_event, instance);
} }
config_handles_close(old); config_handles_close(old);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "configured %d NCD client handles", count); DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "configured %d group client requests", count);
return 0; return 0;
fail: fail:
config_handles_close(updated); config_handles_close(updated);
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "NCD client configuration rejected; previous handles retained"); DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "group client configuration rejected; previous requests retained");
return -1; return -1;
} }

2
src/utun_instance.c

@ -446,7 +446,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
int ch_count = 0; int ch_count = 0;
while (ch) { while (ch) {
struct CONFIG_CONN_HANDLE* next = ch->next; struct CONFIG_CONN_HANDLE* next = ch->next;
node_conn_direct_close(ch->handle); topo_group_peer_close(ch->request);
u_free(ch); ch_count++; u_free(ch); ch_count++;
ch = next; ch = next;
} }

2
src/utun_instance.h

@ -79,7 +79,7 @@ struct conn_queue_entry {
struct CONFIG_CONN_HANDLE { struct CONFIG_CONN_HANDLE {
uint64_t node_id; uint64_t node_id;
char name[MAX_CONN_NAME_LEN]; char name[MAX_CONN_NAME_LEN];
struct NODE_CONN_DIRECT* handle; struct TOPO_PEER_REQUEST* request;
struct CONFIG_CONN_HANDLE* next; struct CONFIG_CONN_HANDLE* next;
}; };

16
tests/bbr_integration/test_bbr_integration.c

@ -108,19 +108,8 @@ static uint64_t now_us(void) {
/* ===== Instance creation ===== */ /* ===== Instance creation ===== */
static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id, static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
const char* priv_hex, const char* pub_hex) { const char* priv_hex, const char* pub_hex) {
struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst));
if (!inst) return NULL;
inst->ua = u;
inst->node_id = node_id;
if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; }
inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool");
inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool");
inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool");
if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; }
inst->connections = queue_new(u, 256, 0, 8, "connections");
if (!inst->connections) { u_free(inst); return NULL; }
struct utun_config* cfg = u_calloc(1, sizeof(*cfg)); struct utun_config* cfg = u_calloc(1, sizeof(*cfg));
if (!cfg) { u_free(inst); return NULL; } if (!cfg) return NULL;
strncpy(cfg->global.my_public_key_hex, pub_hex, MAX_KEY_LEN - 1); strncpy(cfg->global.my_public_key_hex, pub_hex, MAX_KEY_LEN - 1);
strncpy(cfg->global.my_private_key_hex, priv_hex, MAX_KEY_LEN - 1); strncpy(cfg->global.my_private_key_hex, priv_hex, MAX_KEY_LEN - 1);
cfg->global.my_node_id = node_id; cfg->global.my_node_id = node_id;
@ -129,8 +118,7 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id,
cfg->global.keepalive_interval = 500; cfg->global.keepalive_interval = 500;
cfg->global.allowed_keys_allow_all = 1; cfg->global.allowed_keys_allow_all = 1;
cfg->global.bbr_max_cwnd = 100000; cfg->global.bbr_max_cwnd = 100000;
inst->config = cfg; return utun_instance_create_from_config(u, cfg);
return inst;
} }
static int add_server(struct UTUN_INSTANCE* inst, const char* name, int port) { static int add_server(struct UTUN_INSTANCE* inst, const char* name, int port) {

38
tests/test_group_recovery.c

@ -3,6 +3,8 @@
#include "utun_instance.h" #include "utun_instance.h"
#include "topo_group.h" #include "topo_group.h"
#include "topo_recovery.h" #include "topo_recovery.h"
#include "topo_group_connect.h"
#include "topo_node_sqlite.h"
#include "etcp.h" #include "etcp.h"
#include "etcp_api.h" #include "etcp_api.h"
#include "node_conn_direct.h" #include "node_conn_direct.h"
@ -265,9 +267,45 @@ static void sender_backpressure(void) {
poll_events(&f); destroy(&f); poll_events(&f); destroy(&f);
} }
static void persistent_request(void) {
struct fixture f; create(&f);
struct TOPO_PEER_REQUEST* request = NULL;
assert(topo_group_peer_open(f.group, f.ids[0], &request) == 0);
struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle;
up(&f, 0); ready(&f, 0);
f.conns[0]->links_up = 0;
etcp_cbk_fire(f.conns[0], ETCP_CBK_EVENT_DOWN); poll_events(&f);
assert(topo_group_peer_phase(request) == TOPO_PEER_FAILED);
assert(peer(&f, 0)->handle == ownership && topo_group_peer_conn(request) == f.conns[0]);
up(&f, 0); ready(&f, 0);
assert(topo_group_peer_phase(request) == TOPO_PEER_READY && peer(&f, 0)->handle == ownership);
topo_group_peer_close(request); destroy(&f);
}
static void manual_planner_cancel(void) {
struct fixture f; create(&f);
assert(sqlite3_open(":memory:", &f.inst->topo_sqlite_db) == SQLITE_OK);
assert(topo_node_sqlite_init(f.inst->topo_sqlite_db) == 0);
f.group->group_type = TOPO_GROUP_TYPE_CHAT; strcpy(f.group->channel_id, "42");
/* The membership table is the only authorization source. */
uint8_t key[32] = {1}, signature[64] = {1};
assert(topo_node_sqlite_channel_put(f.inst->topo_sqlite_db, "42", "test", 0, key, NULL, key, NULL, signature) == 0);
assert(topo_node_sqlite_member_block_put(f.inst->topo_sqlite_db, "42", f.ids[0], signature, 1, signature, 1,
key, key, "{}", NULL, 0, signature) == 0);
assert(topo_group_connect_node_once(f.group, f.ids[0]) == 0);
struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle;
assert(topo_group_connect_node_once(f.group, f.ids[0]) == 0 && peer(&f, 0)->handle == ownership);
up(&f, 0); poll_events(&f);
assert(topo_group_connect_active_count(f.group) == 0); /* UP does not finish a group attempt. */
topo_groups_remove_group(f.inst->topo_groups, f.group->group_id); f.group = NULL;
poll_events(&f); destroy(&f);
}
int main(void) { int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO); debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
utun_instance_set_tun_init_enabled(0); utun_instance_set_tun_init_enabled(0);
persistent_request();
manual_planner_cancel();
sender_backpressure(); sender_backpressure();
partial_and_external_routes(); partial_and_external_routes();
stalled_and_exhausted(); stalled_and_exhausted();

7
tests/test_ncd_config.c

@ -3,6 +3,7 @@
#include <string.h> #include <string.h>
#include "utun_instance.h" #include "utun_instance.h"
#include "config_parser.h" #include "config_parser.h"
#include "routing_layer/topo_group.h"
#include "transport_layer/etcp.h" #include "transport_layer/etcp.h"
#include "transport_layer/etcp_connections.h" #include "transport_layer/etcp_connections.h"
#include "transport_layer/node_conn_direct.h" #include "transport_layer/node_conn_direct.h"
@ -46,12 +47,12 @@ int main(void) {
first.next = &second; second.next = &third; client.links = &first; first.next = &second; second.next = &third; client.links = &first;
assert(init_client_connections(inst, &client) == 0); assert(init_client_connections(inst, &client) == 0);
assert(inst->config_conn_handles && !inst->config_conn_handles->next); assert(inst->config_conn_handles && !inst->config_conn_handles->next);
struct ETCP_CONN* conn = node_conn_direct_get_conn(inst->config_conn_handles->handle); struct ETCP_CONN* conn = topo_group_peer_conn(inst->config_conn_handles->request);
assert(conn && links(conn) == 3 && !conn->fin_wait); assert(conn && links(conn) == 3 && !conn->fin_wait);
struct NODE_CONN_DIRECT* chat = NULL; struct NODE_CONN_DIRECT* chat = NULL;
assert(node_conn_direct_open(inst, conn->peer_node_id, NULL, NULL, &chat, NULL) == NCD_REUSED); assert(node_conn_direct_open(inst, conn->peer_node_id, NULL, NULL, &chat, NULL) == NCD_REUSED);
assert(init_client_connections(inst, &client) == 0); assert(init_client_connections(inst, &client) == 0);
assert(node_conn_direct_get_conn(inst->config_conn_handles->handle) == conn && links(conn) == 3 && !conn->fin_wait); assert(topo_group_peer_conn(inst->config_conn_handles->request) == conn && links(conn) == 3 && !conn->fin_wait);
/* Ошибка нового набора не освобождает предыдущие handles. */ /* Ошибка нового набора не освобождает предыдущие handles. */
struct CFG_CLIENT invalid = { .name = "invalid", .peer_public_key_hex = "not-a-key" }; struct CFG_CLIENT invalid = { .name = "invalid", .peer_public_key_hex = "not-a-key" };
@ -85,7 +86,7 @@ int main(void) {
} }
third.local_srv = NULL; third.local_srv = NULL;
assert(init_client_connections(inst, &client) == 0); assert(init_client_connections(inst, &client) == 0);
conn = node_conn_direct_get_conn(inst->config_conn_handles->handle); conn = topo_group_peer_conn(inst->config_conn_handles->request);
assert(conn && links(conn) == 1 && conn->links->is_tcp && !conn->links->conn); assert(conn && links(conn) == 1 && conn->links->is_tcp && !conn->links->conn);
assert(conn->links->reality_set && !strcmp(conn->links->reality.server_name, client.reality.server_name)); assert(conn->links->reality_set && !strcmp(conn->links->reality.server_name, client.reality.server_name));
assert(init_client_connections(inst, NULL) == 0); assert(init_client_connections(inst, NULL) == 0);

Loading…
Cancel
Save