diff --git a/src/Makefile.am b/src/Makefile.am index 9b887fee..de068c32 100644 --- a/src/Makefile.am +++ b/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 \ diff --git a/src/routing_layer/conn_mgr.c b/src/routing_layer/conn_mgr.c index 1fd68a55..28b24bc7 100644 --- a/src/routing_layer/conn_mgr.c +++ b/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); diff --git a/src/routing_layer/conn_mgr.h b/src/routing_layer/conn_mgr.h index 0c033d1c..dc2a4bf8 100644 --- a/src/routing_layer/conn_mgr.h +++ b/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-таймаутом ==================== */ /** diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index c99be2bf..b7ada200 100644 --- a/src/routing_layer/topo_group.c +++ b/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); } diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 25068ab7..9d83300f 100644 --- a/src/routing_layer/topo_group.h +++ b/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; // список активных процедур восстановления }; /** diff --git a/src/routing_layer/topo_recovery.c b/src/routing_layer/topo_recovery.c new file mode 100644 index 00000000..b1705a68 --- /dev/null +++ b/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 +#include +#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; +} diff --git a/src/routing_layer/topo_recovery.h b/src/routing_layer/topo_recovery.h new file mode 100644 index 00000000..d387982d --- /dev/null +++ b/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 + * + * Группировка: <мы> -> -> + * Каждый 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 +#include + +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 diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index c356eb94..7f1f9cc2 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/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