Browse Source

topo_recovery: каскадное восстановление узлов с группировкой по next_hop

При разрыве ETCP-соединения модуль собирает каскадно отвалившиеся узлы:
- add_node из hop_list вычисляет next_hop относительно failed_peer,
  группирует узлы в отдельные recovery-контексты (direct/indirect списки)
- start запускает асинхронный перебор для каждой группы параллельно
- is_next_hop-узел пробуется первым с приоритетом по мин. RTT
- таймаут 2с на попытку, conn_mgr_cancel_callback при отмене
- отмена: topo_group_new_conn (on_up) находит узел в любом ctx -> ctx целиком
- завершение: оба списка пусты -> освобождение ctx
topo_upd
Evgeny 2 months ago
parent
commit
4215349349
  1. 2
      src/Makefile.am
  2. 38
      src/routing_layer/conn_mgr.c
  3. 16
      src/routing_layer/conn_mgr.h
  4. 13
      src/routing_layer/topo_group.c
  5. 2
      src/routing_layer/topo_group.h
  6. 265
      src/routing_layer/topo_recovery.c
  7. 82
      src/routing_layer/topo_recovery.h
  8. 1
      tools/chatgui/libutun/CMakeLists.txt

2
src/Makefile.am

@ -16,6 +16,7 @@ utun_CORE_SOURCES = \
routing_layer/topo_node_sqlite.c \
routing_layer/route_connectivity.c \
routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \
db_sync.c \
routing_layer/routing.c \
tun_if.c \
@ -71,6 +72,7 @@ libutun_a_SOURCES = \
routing_layer/topo_node_sqlite.c \
routing_layer/route_connectivity.c \
routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \
db_sync.c \
routing_layer/routing.c \
tun_if.c \

38
src/routing_layer/conn_mgr.c

@ -40,6 +40,7 @@ static void cm_entry_destroy(struct CONN_MGR_ENTRY* entry);
static void cm_exchange_probe_retry_cb(void* arg);
static uint64_t cm_get_node_max_probe_time(struct TOPO_NODEQ* nq);
static int cm_is_rtt_fresh(struct TOPO_NODEQ* nq, uint64_t now_tb);
static void cm_clear_nodeinfo(struct CONN_MGR* mgr, uint64_t node_id);
static void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry);
static void cm_handle_direct_req(struct ETCP_CONN* conn, const uint8_t* data, size_t len);
@ -155,8 +156,16 @@ static struct CONN_MGR_ENTRY* cm_find_entry(struct CONN_MGR* mgr, uint64_t node_
static void cm_deliver_result(struct CONN_MGR_ENTRY* entry, int result) {
struct cm_cb_node* node = entry->cb_list;
while (node) { struct cm_cb_node* next = node->next; node->cb(result, entry->node_id, node->arg); u_free(node); node = next; }
entry->cb_list = NULL;
int delivered = 0;
while (node) {
struct cm_cb_node* next = node->next;
if (!node->cancelled) { node->cb(result, entry->node_id, node->arg); delivered++; }
u_free(node);
node = next;
}
if (delivered == 0 && entry->state == CONN_MGR_STATE_CONNECTING)
cm_clear_nodeinfo(entry->mgr, entry->node_id);
}
static void cm_update_nodeinfo(struct CONN_MGR_ENTRY* entry) {
@ -284,6 +293,33 @@ int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id) {
return CONN_MGR_OK;
}
int conn_mgr_cancel_callback(struct CONN_MGR* mgr, uint64_t node_id,
conn_mgr_connect_callback_t cb, void* cb_arg) {
if (!mgr) return CONN_MGR_ERR_INTERNAL;
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id);
if (!entry) return CONN_MGR_ERR_NOT_FOUND;
int cb_count = 0;
struct cm_cb_node* target = NULL;
for (struct cm_cb_node* n = entry->cb_list; n; n = n->next) {
cb_count++;
if (n->cb == cb && n->arg == cb_arg) target = n;
}
if (!target) return CONN_MGR_ERR_NOT_FOUND;
if (cb_count == 1) {
cm_entry_destroy(entry);
u_free(target);
entry->cb_list = NULL;
cm_clear_nodeinfo(mgr, node_id);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: cancel connect to 0x%016llx (single cb)", (unsigned long long)node_id);
} else {
target->cancelled = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: cancel cb for 0x%016llx (cbs=%d)", (unsigned long long)node_id, cb_count);
}
return CONN_MGR_OK;
}
int conn_mgr_set_idle_timeout(struct CONN_MGR* mgr, uint64_t node_id, uint32_t timeout_ms) {
if (!mgr) return CONN_MGR_ERR_INTERNAL;
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id);

16
src/routing_layer/conn_mgr.h

@ -251,6 +251,7 @@ struct cm_cb_node {
conn_mgr_connect_callback_t cb; ///< Функция-callback
void* arg; ///< Пользовательский аргумент
struct cm_cb_node* next; ///< Следующий callback в списке
uint8_t cancelled; ///< 1 = callback помечен как недействительный
};
/**
@ -371,6 +372,21 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id,
*/
int conn_mgr_disconnect_node(struct CONN_MGR* mgr, uint64_t node_id);
/**
* @brief Отменяет ожидающий callback подключения к ноде.
* @param mgr менеджер
* @param node_id идентификатор ноды
* @param cb callback для отмены
* @param cb_arg аргумент callback для точного сопоставления
* @return CONN_MGR_OK или CONN_MGR_ERR_NOT_FOUND/INTERNAL
*
* Если callback единственный в списке — полностью отменяет попытку подключения
* (уничтожает entry, таймеры, pending-структуры). Если есть другие callback'и —
* только помечает этот callback как cancelled, он будет пропущен при доставке.
*/
int conn_mgr_cancel_callback(struct CONN_MGR* mgr, uint64_t node_id,
conn_mgr_connect_callback_t cb, void* cb_arg);
/* ==================== Управление idle-таймаутом ==================== */
/**

13
src/routing_layer/topo_group.c

@ -23,6 +23,7 @@
#include "route_connectivity.h"
#include "control_server.h"
#include "conn_mgr.h"
#include "topo_recovery.h"
#define TOPO_GROUP_UTUN 0x8000000000000000ULL
@ -236,6 +237,8 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
if (!group) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group_id=%016llx", (unsigned long long)group->group_id);
topo_recovery_cancel_all(group);
struct ll_entry* e;
while ((e = queue_data_get(group->senders_list)) != NULL) queue_entry_free(e);
queue_free(group->senders_list);
@ -395,6 +398,8 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance is NULL"); return; }
if (!conn->instance->rt) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance->rt is NULL"); return; }
topo_recovery_cancel_for_node(group, conn->peer_node_id);
struct TOPO_NODEQ* peer_nq = topo_node_find_by_id(group, conn->peer_node_id);
topo_group_add_to_senders(group, conn);
@ -416,17 +421,19 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
struct ROUTE_TABLE* rt = conn->instance->rt;
int nodes_removed = 0;
int cascaded = 0;
struct ll_entry* node_entry = group->nodes ? group->nodes->head : NULL;
while (node_entry) {
struct ll_entry* next = node_entry->next;
struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)node_entry;
if (topo_group_remove_path(nq, conn) == 1) {
uint64_t key = nq->hash_node_id;
if (key != conn->peer_node_id) { topo_recovery_add_node(group, nq, conn->peer_node_id); cascaded++; }
if (group->group_type != TOPO_GROUP_TYPE_CHAT && rt) route_delete(rt, nq);
if (conn->instance && conn->instance->control_srv) control_server_notify_node_removed(conn->instance->control_srv, nq->hash_node_id);
if (conn->instance && conn->instance->control_srv) control_server_notify_node_removed(conn->instance->control_srv, key);
nq->dirty = 1;
route_connectivity_cancel_node(conn->instance, nq);
if (nq->paths) { queue_free(nq->paths); nq->paths = NULL; }
uint64_t key = nq->hash_node_id;
topo_node_free_lists(group, nq);
struct ll_entry* entry = node_entry;
if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); }
@ -442,6 +449,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
e = group->senders_list->head;
while (e) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; if (item->conn == conn) { queue_remove_data(group->senders_list, e); queue_entry_free(e); break; } e = e->next; }
if (cascaded > 0) topo_recovery_start(group);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d", conn->log_name, nodes_removed);
}

2
src/routing_layer/topo_group.h

@ -42,6 +42,7 @@ extern "C" {
struct UTUN_INSTANCE;
struct NAT_DETECTION;
struct CONN_MGR;
struct TOPO_RECOVERY_CTX;
/** 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,
@ -130,6 +131,7 @@ struct TOPO_GROUP {
uint8_t ed25519_public_key[SC_PUBKEY_SIZE];
char channel_id[64]; // channel_id для групп типа CHAT
struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы
struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления
};
/**

265
src/routing_layer/topo_recovery.c

@ -0,0 +1,265 @@
/**
* @file topo_recovery.c
* @brief Реализация восстановления каскадных узлов.
*
* Цепочка: add_node (группировка по next_hop) → start (запуск всех ctx) → try_next (цикл)
* try_next: сперва is_next_hop-узел группы, затем остальные по мин. RTT
* ├─ recovery_callback (ok/fail): удалить узел → try_next
* └─ timeout_callback (2с): conn_mgr_cancel_callback → удалить → try_next
* Завершение: оба списка пусты → освобождение ctx
* Отмена: cancel_for_node (on_up) → cancel_callback + освобождение ctx
*/
#include <stdlib.h>
#include <string.h>
#include "../lib/platform_compat.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "../lib/u_async.h"
#include "utun_instance.h"
#include "topo_node.h"
#include "topo_group.h"
#include "topo_recovery.h"
#include "conn_mgr.h"
#include "config_parser.h"
#include "../transport_layer/etcp_connections.h"
#define RECOVERY_CAPACITY_INIT 16
/* Выделение и инициализация контекста восстановления */
static struct TOPO_RECOVERY_CTX* topo_recovery_ctx_alloc(struct TOPO_GROUP* group, uint64_t next_hop_id) {
struct TOPO_RECOVERY_CTX* ctx = u_calloc(1, sizeof(struct TOPO_RECOVERY_CTX));
if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery: alloc ctx failed"); return NULL; }
ctx->instance = group->instance; ctx->group = group; ctx->next_hop_id = next_hop_id;
ctx->capacity[0] = RECOVERY_CAPACITY_INIT; ctx->capacity[1] = RECOVERY_CAPACITY_INIT;
ctx->nodes[0] = u_calloc(ctx->capacity[0], sizeof(struct TOPO_RECOVERY_NODE));
ctx->nodes[1] = u_calloc(ctx->capacity[1], sizeof(struct TOPO_RECOVERY_NODE));
if (!ctx->nodes[0] || !ctx->nodes[1]) { u_free(ctx->nodes[0]); u_free(ctx->nodes[1]); u_free(ctx); return NULL; }
return ctx;
}
static void topo_recovery_ctx_free(struct TOPO_RECOVERY_CTX* ctx) {
if (!ctx) return;
u_free(ctx->nodes[0]); u_free(ctx->nodes[1]);
u_free(ctx);
}
/* Аналог cm_has_direct_ip: есть ли адреса с прямым доступом (PUBLIC/UNKNOWN/EIM/DIRECT) */
static int node_has_direct_ip(const struct TOPO_NODEQ* nq) {
if (!nq || !nq->node) return 0;
for (const struct TOPO_ADDR4* a = nq->node->v4_addrs; a; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
for (const struct TOPO_SOCKMETA4* m = nq->node->v4_sock_meta; m; m = m->next) {
if (m->id != a->socket_id) continue;
if (m->config_type == CFG_SERVER_TYPE_PUBLIC || m->config_type == CFG_SERVER_TYPE_UNKNOWN
|| m->nat_type == NAT_VERIFIED_EIM || m->nat_type == NAT_VERIFIED_DIRECT)
return 1;
}
}
return 0;
}
/* Минимальный RTT из данных connectivity-проб (interface/nat/real). 0xFFFF если неизвестен */
static uint16_t node_get_min_rtt(const struct TOPO_NODEQ* nq) {
uint16_t rtt = 0xFFFF;
if (nq->connectivity.interface_status == PROBE_RESULT_REACHABLE && nq->connectivity.interface_min_rtt < rtt)
rtt = nq->connectivity.interface_min_rtt;
if (nq->connectivity.nat_status == PROBE_RESULT_REACHABLE && nq->connectivity.nat_min_rtt < rtt)
rtt = nq->connectivity.nat_min_rtt;
if (nq->connectivity.real_status == PROBE_RESULT_REACHABLE && nq->connectivity.real_min_rtt < rtt)
rtt = nq->connectivity.real_min_rtt;
return rtt;
}
/* Вычисляет next_hop относительно failed_peer из hop_list.
* hop_list = [...узлы пути...], последний элемент = failed_peer (форвардил NODEINFO нам).
* next_hop = элемент перед failed_peer (hop_list[hop_count-2]).
* Если failed_peer не найден в конце — fallback: node_id сам себе next_hop. */
static uint64_t find_next_hop(const struct TOPO_NODEQ* nq, uint64_t failed_peer) {
if (nq->hop_count >= 2 && nq->hop_list && nq->hop_list[nq->hop_count - 1] == failed_peer)
return nq->hop_list[nq->hop_count - 2];
return nq->hash_node_id;
}
/* Ищет незапущенный (started==0) контекст по next_hop_id */
static struct TOPO_RECOVERY_CTX* topo_recovery_find_by_next_hop(struct TOPO_GROUP* group, uint64_t next_hop_id) {
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list;
while (ctx) { if (!ctx->started && ctx->next_hop_id == next_hop_id) return ctx; ctx = ctx->next; }
return NULL;
}
/* Вынимает контекст из связного списка group->recovery_list */
static void topo_recovery_ctx_remove(struct TOPO_GROUP* group, struct TOPO_RECOVERY_CTX* ctx) {
struct TOPO_RECOVERY_CTX** pp = &group->recovery_list;
while (*pp) { if (*pp == ctx) { *pp = ctx->next; return; } pp = &(*pp)->next; }
}
static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx);
/* Коллбэк conn_mgr: подключение удалось или провалилось */
static void topo_recovery_callback(int result, uint64_t node_id, void* arg) {
struct TOPO_RECOVERY_CTX* ctx = (struct TOPO_RECOVERY_CTX*)arg;
if (node_id != ctx->current_node_id) return;
if (ctx->connect_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->connect_timer); ctx->connect_timer = NULL; }
int list = ctx->current_list;
for (size_t i = 0; i < ctx->count[list]; i++)
if (ctx->nodes[list][i].node_id == node_id) { ctx->nodes[list][i] = ctx->nodes[list][ctx->count[list] - 1]; ctx->count[list]--; break; }
if (result == CONN_MGR_OK)
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx reconnected to %016llx", (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id);
else
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx connect to %016llx failed: %d", (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id, result);
ctx->current_node_id = 0;
topo_recovery_try_next(ctx);
}
/* Таймаут 2с: отменяет conn_mgr-коллбэк, удаляет узел, переходит к следующему */
static void topo_recovery_timeout_cb(void* arg) {
struct TOPO_RECOVERY_CTX* ctx = (struct TOPO_RECOVERY_CTX*)arg;
ctx->connect_timer = NULL;
uint64_t node_id = ctx->current_node_id;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx timeout connecting to %016llx", (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id);
conn_mgr_cancel_callback(ctx->group->conn_mgr, node_id, topo_recovery_callback, ctx);
int list = ctx->current_list;
for (size_t i = 0; i < ctx->count[list]; i++)
if (ctx->nodes[list][i].node_id == node_id) { ctx->nodes[list][i] = ctx->nodes[list][ctx->count[list] - 1]; ctx->count[list]--; break; }
ctx->current_node_id = 0;
topo_recovery_try_next(ctx);
}
/* Выбирает узел: сперва is_next_hop с мин. RTT, за ним — остальные по мин. RTT */
static struct TOPO_RECOVERY_NODE* topo_recovery_find_best(struct TOPO_RECOVERY_CTX* ctx, int list) {
struct TOPO_RECOVERY_NODE* best = NULL;
for (size_t i = 0; i < ctx->count[list]; i++)
if (ctx->nodes[list][i].is_next_hop && (!best || ctx->nodes[list][i].min_rtt < best->min_rtt))
best = &ctx->nodes[list][i];
if (best) return best;
for (size_t i = 0; i < ctx->count[list]; i++)
if (!best || ctx->nodes[list][i].min_rtt < best->min_rtt)
best = &ctx->nodes[list][i];
return best;
}
static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx) {
while (1) {
if (ctx->count[ctx->current_list] == 0) {
ctx->current_list++;
if (ctx->current_list > 1) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx all nodes exhausted, no reconnections", (unsigned long long)ctx->next_hop_id);
topo_recovery_ctx_remove(ctx->group, ctx);
topo_recovery_ctx_free(ctx);
return;
}
continue;
}
struct TOPO_RECOVERY_NODE* best = topo_recovery_find_best(ctx, ctx->current_list);
uint64_t node_id = best->node_id;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx trying %016llx rtt=%u list=%s %s",
(unsigned long long)ctx->next_hop_id, (unsigned long long)node_id, best->min_rtt,
ctx->current_list == 0 ? "direct" : "indirect", best->is_next_hop ? "[next_hop]" : "");
ctx->current_node_id = node_id;
int rc = conn_mgr_connect_node(ctx->group->conn_mgr, node_id, 0, topo_recovery_callback, ctx);
if (rc != CONN_MGR_OK) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: conn_mgr_connect_node %016llx returned %d", (unsigned long long)node_id, rc);
*best = ctx->nodes[ctx->current_list][ctx->count[ctx->current_list] - 1]; ctx->count[ctx->current_list]--;
ctx->current_node_id = 0;
continue;
}
ctx->connect_timer = uasync_set_timeout(ctx->instance->ua, TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10,
ctx, topo_recovery_timeout_cb, "recovery_timeout");
return;
}
}
void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_NODEQ* nq, uint64_t failed_peer) {
if (!group || !nq || !nq->node) return;
uint64_t node_id = nq->hash_node_id;
uint64_t next_hop = find_next_hop(nq, failed_peer);
struct TOPO_RECOVERY_CTX* ctx = topo_recovery_find_by_next_hop(group, next_hop);
if (!ctx) {
ctx = topo_recovery_ctx_alloc(group, next_hop);
if (!ctx) return;
ctx->next = group->recovery_list;
group->recovery_list = ctx;
}
int list = node_has_direct_ip(nq) ? 0 : 1;
uint16_t rtt = node_get_min_rtt(nq);
if (ctx->count[list] >= ctx->capacity[list]) {
size_t new_cap = ctx->capacity[list] * 2;
struct TOPO_RECOVERY_NODE* new_nodes = u_realloc(ctx->nodes[list], new_cap * sizeof(struct TOPO_RECOVERY_NODE));
if (!new_nodes) return;
ctx->nodes[list] = new_nodes; ctx->capacity[list] = new_cap;
}
ctx->nodes[list][ctx->count[list]].node_id = node_id;
ctx->nodes[list][ctx->count[list]].min_rtt = rtt;
ctx->nodes[list][ctx->count[list]].is_next_hop = (uint8_t)(node_id == next_hop ? 1 : 0);
ctx->count[list]++;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "recovery: next=%016llx add node %016llx list=%s rtt=%u is_nh=%d count=%zu",
(unsigned long long)next_hop, (unsigned long long)node_id,
list == 0 ? "direct" : "indirect", rtt, node_id == next_hop ? 1 : 0, ctx->count[list]);
}
void topo_recovery_start(struct TOPO_GROUP* group) {
if (!group) return;
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list;
int started_ctxs = 0;
while (ctx) {
struct TOPO_RECOVERY_CTX* next_ctx = ctx->next;
if (!ctx->started) {
if (ctx->count[0] + ctx->count[1] == 0) {
topo_recovery_ctx_remove(group, ctx);
topo_recovery_ctx_free(ctx);
} else {
ctx->started = 1; ctx->current_list = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: start next=%016llx direct=%zu indirect=%zu",
(unsigned long long)ctx->next_hop_id, ctx->count[0], ctx->count[1]);
topo_recovery_try_next(ctx);
started_ctxs++;
}
}
ctx = next_ctx;
}
if (started_ctxs == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "recovery: start called but no pending ctx found");
}
}
void topo_recovery_cancel_for_node(struct TOPO_GROUP* group, uint64_t node_id) {
if (!group) return;
struct TOPO_RECOVERY_CTX** pp = &group->recovery_list;
while (*pp) {
struct TOPO_RECOVERY_CTX* ctx = *pp;
int found = 0;
for (int list = 0; list < 2 && !found; list++)
for (size_t i = 0; i < ctx->count[list] && !found; i++)
if (ctx->nodes[list][i].node_id == node_id) found = 1;
if (found) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: cancel next=%016llx (node %016llx came up)",
(unsigned long long)ctx->next_hop_id, (unsigned long long)node_id);
if (ctx->current_node_id != 0)
conn_mgr_cancel_callback(group->conn_mgr, ctx->current_node_id, topo_recovery_callback, ctx);
if (ctx->connect_timer) { uasync_cancel_timeout(group->instance->ua, ctx->connect_timer); ctx->connect_timer = NULL; }
*pp = ctx->next;
topo_recovery_ctx_free(ctx);
} else {
pp = &ctx->next;
}
}
}
void topo_recovery_cancel_all(struct TOPO_GROUP* group) {
if (!group) return;
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list;
while (ctx) {
struct TOPO_RECOVERY_CTX* next = ctx->next;
if (ctx->current_node_id != 0)
conn_mgr_cancel_callback(group->conn_mgr, ctx->current_node_id, topo_recovery_callback, ctx);
if (ctx->connect_timer) { uasync_cancel_timeout(group->instance->ua, ctx->connect_timer); ctx->connect_timer = NULL; }
topo_recovery_ctx_free(ctx);
ctx = next;
}
group->recovery_list = NULL;
}

82
src/routing_layer/topo_recovery.h

@ -0,0 +1,82 @@
/**
* @file topo_recovery.h
* @brief Восстановление каскадно отвалившихся узлов после разрыва ETCP-соединения.
*
* Когда рвётся соединение, узлы достижимые только через него становятся недоступны.
* Модуль собирает их перед удалением из BGP-таблицы и пробует переподключиться:
* 1. add_node — из hop_list вычисляет next_hop относительно failed_peer,
* группирует узлы по next_hop в отдельные recovery-контексты
* 2. Классификация на direct (есть прямые IP) / indirect (нет)
* 3. start — для каждого контекста запускает асинхронный перебор:
* сперва пробует next_hop-узел (is_next_hop=1), затем остальные по мин. RTT
* 4. Таймаут 2с на каждую попытку conn_mgr_connect_node
*
* Группировка: <мы> -> <failed_peer N> -> <next_hop A, B...>
* Каждый next_hop со своим subtree — отдельный recovery-контекст.
* Восстановление next_hop автоматически оживляет его subtree через BGP.
*
* Отмена: topo_group_new_conn (узел появился в сети) сканирует все active recovery,
* при совпадении отменяет контекст целиком. Самоуничтожение при исчерпании списков.
*/
#ifndef TOPO_RECOVERY_H
#define TOPO_RECOVERY_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
#include <stddef.h>
struct TOPO_GROUP;
struct TOPO_NODEQ;
#define TOPO_RECOVERY_CONNECT_TIMEOUT_MS 2000
struct TOPO_RECOVERY_NODE {
uint64_t node_id;
uint16_t min_rtt; /* 0.1ms, из connectivity-проб (interface/nat/real) */
uint8_t is_next_hop; /* 1 = прямой downstream отвалившегося узла, пробуется первым */
};
struct TOPO_RECOVERY_CTX {
struct TOPO_RECOVERY_CTX* next; /* следующий в group->recovery_list */
struct UTUN_INSTANCE* instance;
struct TOPO_GROUP* group;
uint64_t next_hop_id; /* ключ группировки: node_id next-hop'а от failed_peer */
struct TOPO_RECOVERY_NODE* nodes[2]; /* [0]=с прямыми IP, [1]=без прямых IP */
size_t count[2]; /* текущее количество в каждом списке */
size_t capacity[2]; /* выделенная ёмкость */
int current_list; /* 0=direct, 1=indirect */
uint64_t current_node_id; /* node_id в текущей попытке, 0=нет активной */
void* connect_timer; /* внешний таймер 2с (uasync) */
uint8_t started; /* 0=сбор узлов (add_node), 1=перебор запущен */
};
/**
* Добавляет узел в pending-контекст, сгруппированный по next_hop относительно failed_peer.
* Вызывается из topo_group_remove_conn ДО topo_node_free_lists (нужен hop_list).
* @param failed_peer node_id отвалившегося пира (conn->peer_node_id)
*/
void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_NODEQ* nq, uint64_t failed_peer);
/**
* Запускает перебор узлов для всех pending-контекстов.
* Вызывается после цикла в topo_group_remove_conn если есть каскадные узлы.
*/
void topo_recovery_start(struct TOPO_GROUP* group);
/**
* Сканирует все active recovery: если node_id найден в любом контексте —
* отменяет текущую попытку connect и освобождает контекст.
* Вызывается из topo_group_new_conn (узел появился в сети).
*/
void topo_recovery_cancel_for_node(struct TOPO_GROUP* group, uint64_t node_id);
/** Отменяет все recovery-контексты. Вызывается из topo_group_destroy. */
void topo_recovery_cancel_all(struct TOPO_GROUP* group);
#ifdef __cplusplus
}
#endif
#endif

1
tools/chatgui/libutun/CMakeLists.txt

@ -66,6 +66,7 @@ set(UTUN_COMMON_SOURCES
${SRC_DIR}/routing_layer/route_connectivity.c
${SRC_DIR}/db_sync.c
${SRC_DIR}/routing_layer/conn_mgr.c
${SRC_DIR}/routing_layer/topo_recovery.c
${TRANSPORT_DIR}/chat_sync.c
${TRANSPORT_DIR}/merkle_sync.c
${TRANSPORT_DIR}/member_sync.c

Loading…
Cancel
Save