diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index bdf39816..3394a098 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -37,20 +37,121 @@ // Вспомогательные функции // ============================================================================ +struct TOPO_PEER_REQUEST { + struct TOPO_GROUP_CONN_ITEM* peer; + struct TOPO_PEER_REQUEST* next; +}; + +static int topo_group_begin_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, int retain); +static int topo_group_peer_allowed(struct TOPO_GROUP* group, uint64_t node_id); + static struct TOPO_GROUP_CONN_ITEM* topo_group_peer(const struct TOPO_GROUP* group, uint64_t peer_id) { for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - if (item->conn && item->conn->peer_node_id == peer_id) return item; + if (item->node_id == peer_id) return item; } return NULL; } int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id) { const struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, peer_id); - return peer && peer->conn->links_up && peer->table_received && peer->table_sent; + return peer && peer->conn && peer->conn->links_up && !peer->conn->close_requested && peer->table_received && peer->table_sent; +} + +static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { + while (peer->requests) { + struct TOPO_PEER_REQUEST* request = peer->requests; + peer->requests = request->next; request->peer = NULL; request->next = NULL; + } + if (peer->handle) node_conn_direct_close(peer->handle); + topo_recovery_changed(peer->group); +} + +static void topo_group_pending_remove(struct TOPO_GROUP_CONN_ITEM* peer) { + struct ll_queue* queue = peer->group->senders_list; + for (struct ll_entry* e = queue->head; e; e = e->next) { + if ((void*)e->data != peer) continue; + queue_remove_data(queue, e); topo_group_peer_free(peer); queue_entry_free(e); return; + } +} + +static void topo_group_peer_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) { + struct TOPO_GROUP_CONN_ITEM* peer = arg; + if (event == NCD_EVENT_UP) { + struct ETCP_CONN* conn = node_conn_direct_get_conn(handle); + if (conn && topo_group_begin_conn(peer->group, conn, 0) == 0) return; + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group join failed peer=%016llx group=%016llx", + (unsigned long long)peer->node_id, (unsigned long long)peer->group->group_id); + } + if (event == NCD_EVENT_TIMEOUT) + 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); + if (peer->conn) topo_group_remove_conn(peer->group, peer->conn, TOPO_REMOVE_TRANSPORT_DOWN); + else topo_group_pending_remove(peer); +} + +int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out) { + if (out) *out = NULL; + if (!group || !out || !node_id || !group->senders_list || group->stopping || !topo_group_peer_allowed(group, node_id)) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot request group peer=%016llx: invalid arguments, stopped group or unknown member", + (unsigned long long)node_id); return -1; + } + struct TOPO_PEER_REQUEST* request = u_calloc(1, sizeof(*request)); + if (!request) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group request allocation failed"); return -1; } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id); + if (!peer) { + struct ll_entry* e = queue_entry_new(sizeof(*peer)); + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group peer allocation failed"); u_free(request); return -1; } + peer = (struct TOPO_GROUP_CONN_ITEM*)e->data; + peer->group = group; peer->node_id = node_id; peer->progress = get_time_tb(); + if (node_conn_direct_open(group->instance, node_id, topo_group_peer_event, peer, &peer->handle, NULL) == NCD_ERR) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group NCD open failed peer=%016llx group=%016llx", + (unsigned long long)node_id, (unsigned long long)group->group_id); + queue_entry_free(e); u_free(request); return -1; + } + queue_data_put(group->senders_list, e); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group peer CONNECTING peer=%016llx group=%016llx", + (unsigned long long)node_id, (unsigned long long)group->group_id); + } + request->peer = peer; request->next = peer->requests; peer->requests = request; *out = request; + return 0; +} + +enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request) { + struct TOPO_GROUP_CONN_ITEM* peer = request ? request->peer : NULL; + if (!peer) return TOPO_PEER_FAILED; + 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; +} + +uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request) { + return request && request->peer ? request->peer->progress : 0; +} + +void topo_group_peer_close(struct TOPO_PEER_REQUEST* request) { + if (!request) return; + struct TOPO_GROUP_CONN_ITEM* peer = request->peer; + if (peer) { + struct TOPO_PEER_REQUEST** p = &peer->requests; + while (*p && *p != request) p = &(*p)->next; + if (*p) *p = request->next; + if (!peer->requests && !peer->retained) { + if (peer->conn) topo_group_leave_peer(peer->group, peer->handle); + else topo_group_pending_remove(peer); + } + } + u_free(request); +} + +static void topo_group_peer_progressed(struct TOPO_GROUP* group, uint64_t node_id) { + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id); + if (peer) peer->progress = get_time_tb(); + topo_recovery_changed(group); } 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; + 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", (unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id, (unsigned long long)peer->exchange_id, peer->table_sent, peer->table_received, @@ -392,6 +493,7 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)se->data; if (peer->conn != conn) continue; peer->exchange_id = 0; peer->table_received = 0; peer->table_sent = 0; + topo_group_peer_progressed(g, peer->node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP,"BGP session resync: peer=%016llx group=%016llx", (unsigned long long)conn->peer_node_id,(unsigned long long)g->group_id); topo_group_send_join_group(g,conn); topo_group_send_table_request(g,conn); @@ -411,7 +513,7 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg for (int i = 0; i < chn; i++) { struct TOPO_GROUP* cg = topo_groups_find(groups, chs[i]); if (cg && cg->group_type == TOPO_GROUP_TYPE_CHAT) - topo_group_new_conn(cg, conn); + topo_group_begin_conn(cg, conn, 0); } } u_free(chs); @@ -460,8 +562,9 @@ static struct TOPO_GROUP* topo_group_create(struct UTUN_INSTANCE* instance, uint /* Полностью разрушает группу: recovery, connect, broadcast, conn_mgr, узлы. */ static void topo_group_destroy(struct TOPO_GROUP* group) { if (!group) return; + group->stopping = 1; DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 0 enter grp=%p id=0x%016llx rec_list=%p connect=%p", - group, (unsigned long long)group->group_id, group->recovery_list, group->connect); + group, (unsigned long long)group->group_id, group->recovery, group->connect); topo_recovery_cancel_all(group); DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 1 recovery_cancel done grp=%p", group); @@ -479,7 +582,7 @@ static void topo_group_destroy(struct TOPO_GROUP* group) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "group destroy: release peer=%016llx group=%016llx", (unsigned long long)node_conn_direct_node_id(item->handle), (unsigned long long)group->group_id); - if (item->handle) { node_conn_direct_close(item->handle); item->handle = NULL; } + topo_group_peer_free(item); queue_entry_free(e); } queue_free(group->senders_list); @@ -778,8 +881,9 @@ void topo_groups_set_node_updated_cb(struct TOPO_GROUPS* groups, topo_node_updat } /* Добавляет пира в группу при UP: дедуп, добавление в senders, запрос таблицы. */ -int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { +static int topo_group_begin_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, int retain) { if (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; } + if (group->stopping) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "join on stopped group"); return -1; } if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance is NULL"); return -1; } if (!conn->instance->rt) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance->rt is NULL"); return -1; } if (!topo_group_peer_allowed(group, conn->peer_node_id)) { @@ -787,6 +891,8 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id); return -1; } + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (peer && retain) peer->retained = 1; /* дедуп: тот же conn стреляет ETCP_CONN_STATUS_UP дважды (UDP-линк, затем TCP-линк). * Если conn уже в senders_list — повторно не обрабатываем, иначе active_conn_count @@ -813,7 +919,7 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { (void*)conn, (unsigned long long)conn->peer_node_id, conn->state, (unsigned long long)group->group_id, group->channel_id); if (topo_group_add_to_senders(group, conn) != 0) return -1; - topo_recovery_cancel_for_node(group, conn->peer_node_id); + topo_group_peer_progressed(group, conn->peer_node_id); DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id); @@ -823,10 +929,20 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { return 0; } +int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + return topo_group_begin_conn(group, conn, 1); +} + /* Обрабатывает DOWN пира: чистит пути, каскадно удаляет узлы, шлёт withdraw. */ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason) { if (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; } if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid instance"); return; } + struct TOPO_GROUP_CONN_ITEM* pending = topo_group_peer(group, conn->peer_node_id); + 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", + (unsigned long long)pending->node_id, (unsigned long long)group->group_id, reason); + topo_group_pending_remove(pending); return; + } bool found_in_list = false; struct ll_entry* e = group->senders_list->head; @@ -846,12 +962,25 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)node_entry; { uint64_t next_hop = nq->node_id; - struct ll_entry* pe = nq->paths ? nq->paths->head : NULL; - while (pe) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)pe; if (p->conn == conn && p->hop_count >= 2) { uint64_t* hop = (uint64_t*)((uint8_t*)p + sizeof(struct TOPO_NODEPATH)); if (hop[p->hop_count - 1] == conn->peer_node_id) next_hop = hop[p->hop_count - 2]; break; } pe = pe->next; } + uint16_t lost_rtt = UINT16_MAX; + for (struct ll_entry* pe = nq->paths ? nq->paths->head : NULL; pe; pe = pe->next) { + struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe; + if (path->conn != conn) continue; + /* links_up уже сброшен; сохраняем последний RTT до удаления пути. */ + if (conn->rtt_last) { + uint32_t total = (uint32_t)conn->rtt_last + path->cumulative_rtt; + lost_rtt = total > UINT16_MAX ? UINT16_MAX : (uint16_t)total; + } + if (path->hop_count >= 2) { + uint64_t* hop = (uint64_t*)((uint8_t*)path + sizeof(*path)); + if (hop[path->hop_count - 1] == conn->peer_node_id) next_hop = hop[path->hop_count - 2]; + } + break; + } if (topo_group_remove_path(nq, conn) == 1) { uint64_t key = nq->node_id; if (reason == TOPO_REMOVE_TRANSPORT_DOWN && key != conn->peer_node_id) { - topo_recovery_add_node(group, nq, next_hop); cascaded++; + topo_recovery_add_node(group, key, next_hop, lost_rtt); cascaded++; } nq->conn_presence = 0; nq->conn_up = 0; topo_fire_nodeinfo_cbk(conn->instance, group, nq); @@ -886,17 +1015,12 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en while (e) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; if (item->conn == conn) { - struct NODE_CONN_DIRECT* h = item->handle; - item->handle = NULL; 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", (void*)conn, (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->channel_id); - if (h) { - node_conn_direct_close(h); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "topo_group_remove_conn: BGP handle released node=0x%016llx", (unsigned long long)conn->peer_node_id); - } break; } e = e->next; @@ -1184,6 +1308,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from } /* Создание линков может вызвать транспортные callbacks: структура группы уже обновлена. */ + topo_group_peer_progressed(group, from->peer_node_id); node_conn_direct_update_node(group->instance, node_id); return 0; @@ -1263,12 +1388,21 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n /* Добавляет conn в senders_list (дедуп), если его там ещё нет. */ static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { if (!group || !conn || !group->senders_list) return -1; + struct TOPO_GROUP_CONN_ITEM* pending = topo_group_peer(group, conn->peer_node_id); + if (pending) { + if (pending->conn && pending->conn != conn) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group peer has conflicting transport peer=%016llx", (unsigned long long)conn->peer_node_id); + return -1; + } + pending->conn = conn; return 0; + } for (struct ll_entry* e = group->senders_list->head; e; e = e->next) if (((struct TOPO_GROUP_CONN_ITEM*)e->data)->conn == conn) return 0; struct ll_entry* e = queue_entry_new(sizeof(struct TOPO_GROUP_CONN_ITEM)); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "senders_add: allocation failed"); return -1; } struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - item->conn = conn; item->handle = NULL; + item->group = group; item->node_id = conn->peer_node_id; + item->conn = conn; item->handle = NULL; item->retained = 1; item->progress = get_time_tb(); if (node_conn_direct_open(conn->instance, conn->peer_node_id, NULL, NULL, &item->handle, NULL) == NCD_ERR) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "senders_add: cannot acquire NCD handle peer=%016llx", (unsigned long long)conn->peer_node_id); queue_entry_free(e); return -1; @@ -1417,7 +1551,7 @@ static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_C DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_join_group: from %s grp=%016llx type=%d", conn->log_name, (unsigned long long)group->group_id, group->group_type); /* Проверка таблицы участников общая для входящих и локальных запросов. */ int existing = topo_group_has_sender(group, conn); - if (topo_group_new_conn(group, conn) == 0 && existing) + if (topo_group_begin_conn(group, conn, 0) == 0 && existing) topo_group_send_table_request(group, conn); } diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index c799c836..9bef7a5d 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -50,6 +50,7 @@ struct NAT_DETECTION; struct CONN_MGR; struct TOPO_RECOVERY_CTX; struct TOPO_GROUP_CONNECT; +struct TOPO_PEER_REQUEST; struct broadcast_ctx; struct radio_ctx; @@ -180,8 +181,13 @@ struct TOPOMSG_ERR_GROUP_MISMATCH { } __attribute__((packed)); struct TOPO_GROUP_CONN_ITEM { - struct ETCP_CONN* conn; + struct TOPO_GROUP* group; + uint64_t node_id; + struct ETCP_CONN* conn; // NULL до UP, в состоянии CONNECTING struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) + struct TOPO_PEER_REQUEST* requests; + uint64_t progress; // время последнего прогресса обмена (timebase) + uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint64_t exchange_id; // текущий исходящий REQUEST_TABLE uint8_t table_received; // TABLE_COMPLETE для exchange_id принят uint8_t table_sent; // ответ на запрос пира целиком передан транспорту @@ -203,18 +209,19 @@ struct TOPO_GROUP { uint64_t group_id; // уникальный идентификатор группы uint8_t group_type; // TOPO_GROUP_TYPE_* struct UTUN_INSTANCE* instance; - struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM, без хеш-индекса + struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM; включает CONNECTING с conn=NULL struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B) struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе 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; // список активных процедур восстановления + struct TOPO_RECOVERY_CTX* recovery; // один последовательный recovery на группу struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT) struct broadcast_ctx* broadcast; // broadcast protocol per-group context struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO) struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов + uint8_t stopping; }; /* Отложенный REQUEST_TABLE: пир запросил таблицу группы, которой у нас ещё нет. @@ -311,6 +318,20 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn); * Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */ int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id); +enum topo_peer_phase { TOPO_PEER_FAILED, TOPO_PEER_CONNECTING, TOPO_PEER_SYNCING, TOPO_PEER_READY }; +/* Запрос присоединения делит попытку с другими запросами этой пары (группа, пир). + * Единственный NCD handle принадлежит группе с начала CONNECTING. REQUEST не + * владеет транспортом. Проверка CHAT-членства общая с входящим JOIN_GROUP. + * close отменяет только свой запрос; последний запрос отменяет незавершённую + * попытку, если её не удерживает другой инициатор. READY сохраняется в группе. + * Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан + * закрыть его. Все операции выполняются в потоке единственного uasync. + * progress — timebase последнего принятого NODEINFO/шага обмена, не polling time. */ +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); +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); + /** * @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable. * diff --git a/src/routing_layer/topo_recovery.c b/src/routing_layer/topo_recovery.c index b6d8d65d..8c3c0d4c 100644 --- a/src/routing_layer/topo_recovery.c +++ b/src/routing_layer/topo_recovery.c @@ -1,16 +1,4 @@ -/** - * @file topo_recovery.c - * @brief Реализация восстановления каскадных узлов (прямые подключения, ncd). - * - * Цепочка: add_node (группировка по next_hop) → start (запуск всех ctx) → try_next (цикл) - * try_next: сперва is_next_hop-узел группы, затем остальные по мин. RTT - * ├─ recovery_callback (ok/fail): удалить узел → try_next - * └─ timeout_callback (2с): закрыть handle → удалить → try_next - * Завершение: список пуст → освобождение ctx - * Отмена: cancel_for_node (on_up) → закрыть handle + освобождение ctx - */ -#include -#include +#include #include "../lib/platform_compat.h" #include "../lib/debug_config.h" #include "../lib/mem.h" @@ -19,232 +7,188 @@ #include "topo_node.h" #include "topo_group.h" #include "topo_recovery.h" -#include "../transport_layer/node_conn_direct.h" #include "../transport_layer/etcp.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 = RECOVERY_CAPACITY_INIT; - ctx->nodes = u_calloc(ctx->capacity, sizeof(struct TOPO_RECOVERY_NODE)); - if (!ctx->nodes) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery: nodes allocation failed"); u_free(ctx); return NULL; } - return ctx; -} - -static void topo_recovery_remove_candidate(struct TOPO_RECOVERY_CTX* ctx, size_t index) { - uint64_t id = ctx->nodes[index].node_id; - ctx->nodes[index] = ctx->nodes[--ctx->count]; - topo_node_registry_unref(ctx->instance->topo_groups, id); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "recovery: release candidate=%016llx remaining=%zu", (unsigned long long)id, ctx->count); -} - -static void topo_recovery_ctx_free(struct TOPO_RECOVERY_CTX* ctx) { - if (!ctx) return; - while (ctx->count) topo_recovery_remove_candidate(ctx, ctx->count - 1); - u_free(ctx->nodes); - u_free(ctx); +struct recovery_node { + uint64_t node_id; + uint16_t rtt; + uint8_t priority, tried, restored; +}; + +struct TOPO_RECOVERY_CTX { + struct TOPO_GROUP* group; + struct recovery_node* nodes; + size_t count, capacity; + struct TOPO_PEER_REQUEST* request; + uint64_t current_node, started_at; + enum topo_peer_phase phase; + void* wake; + void* timer; + unsigned attempts; + uint8_t started; +}; + +static void recovery_step(void* arg); + +static void recovery_finish(struct TOPO_RECOVERY_CTX* ctx) { + struct TOPO_GROUP* group = ctx->group; + group->recovery = NULL; + if (ctx->wake) uasync_call_soon_cancel(group->instance->ua, ctx->wake); + if (ctx->timer) uasync_cancel_timeout(group->instance->ua, ctx->timer); + topo_group_peer_close(ctx->request); + for (size_t i = 0; i < ctx->count; i++) topo_node_registry_unref(group->instance->topo_groups, ctx->nodes[i].node_id); + u_free(ctx->nodes); u_free(ctx); } -/* Ищет незапущенный (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 int recovery_has_route(struct TOPO_GROUP* group, uint64_t node_id) { + struct TOPO_GROUP_NODE* node = topo_node_find_by_id(group, node_id); + for (struct ll_entry* e = node && node->paths ? node->paths->head : NULL; e; e = e->next) { + struct ETCP_CONN* conn = ((struct TOPO_NODEPATH*)e)->conn; + if (conn && conn->links_up && !conn->close_requested && topo_group_peer_ready(group, conn->peer_node_id)) return 1; + } + return 0; } -static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx); - -/* Коллбэк ncd: подключение удалось или провалилось */ -static void topo_recovery_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { - struct TOPO_RECOVERY_CTX* ctx = (struct TOPO_RECOVERY_CTX*)arg; - struct ETCP_CONN* conn = node_conn_direct_get_conn(h); - uint64_t node_id = node_conn_direct_node_id(h); - 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; } - ctx->current_node_id = 0; - ctx->current_handle = NULL; - if (event == NCD_EVENT_UP && conn) { - // Передача владения может отменить другие recovery. Свой ctx исключаем из списка заранее. - topo_recovery_ctx_remove(ctx->group, ctx); - if (topo_group_new_conn(ctx->group, conn) == 0) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: reconnected node=%016llx", (unsigned long long)node_id); - node_conn_direct_close(h); - topo_recovery_ctx_free(ctx); - return; +static size_t recovery_remaining(struct TOPO_RECOVERY_CTX* ctx) { + size_t remaining = 0; + for (size_t i = 0; i < ctx->count; i++) { + struct recovery_node* node = &ctx->nodes[i]; + int restored = recovery_has_route(ctx->group, node->node_id); + if (restored != node->restored) { + node->restored = (uint8_t)restored; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery route: group=%016llx node=%016llx restored=%d", + (unsigned long long)ctx->group->group_id, (unsigned long long)node->node_id, restored); } - DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: BGP ownership transfer failed node=%016llx", (unsigned long long)node_id); - ctx->next = ctx->group->recovery_list; ctx->group->recovery_list = ctx; + if (!restored) remaining++; } - node_conn_direct_close(h); - for (size_t i = 0; i < ctx->count; i++) - if (ctx->nodes[i].node_id == node_id) { topo_recovery_remove_candidate(ctx, i); break; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: attempt finished node=%016llx event=%d", (unsigned long long)node_id, event); - topo_recovery_try_next(ctx); + return remaining; } -/* Таймаут 2с: закрывает handle, удаляет узел, переходит к следующему */ -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); - if (ctx->current_handle) { - node_conn_direct_close(ctx->current_handle); - ctx->current_handle = NULL; +static struct recovery_node* recovery_best(struct TOPO_RECOVERY_CTX* ctx) { + struct recovery_node* best = NULL; + for (size_t i = 0; i < ctx->count; i++) { + struct recovery_node* node = &ctx->nodes[i]; + if (node->tried || node->restored) continue; + if (!best || node->priority > best->priority || (node->priority == best->priority && + (node->rtt < best->rtt || (node->rtt == best->rtt && node->node_id < best->node_id)))) best = node; } - for (size_t i = 0; i < ctx->count; i++) - if (ctx->nodes[i].node_id == node_id) { topo_recovery_remove_candidate(ctx, i); break; } - ctx->current_node_id = 0; - topo_recovery_try_next(ctx); + return best; } -/* Выбирает узел: сперва is_next_hop с мин. RTT, за ним — остальные по мин. RTT */ -static struct TOPO_RECOVERY_NODE* topo_recovery_find_best(struct TOPO_RECOVERY_CTX* ctx) { - struct TOPO_RECOVERY_NODE* best = NULL; - for (size_t i = 0; i < ctx->count; i++) - if (ctx->nodes[i].is_next_hop && (!best || ctx->nodes[i].min_rtt < best->min_rtt)) - best = &ctx->nodes[i]; - if (best) return best; - for (size_t i = 0; i < ctx->count; i++) - if (!best || ctx->nodes[i].min_rtt < best->min_rtt) - best = &ctx->nodes[i]; - return best; +static void recovery_timeout(void* arg) { + struct TOPO_RECOVERY_CTX* ctx = arg; + ctx->timer = NULL; + recovery_step(ctx); +} + +static void recovery_wake(void* arg) { + struct TOPO_RECOVERY_CTX* ctx = arg; + ctx->wake = NULL; + recovery_step(ctx); } -static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx) { - while (1) { - if (ctx->count == 0) { - 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; +static void recovery_step(void* arg) { + struct TOPO_RECOVERY_CTX* ctx = arg; + struct TOPO_GROUP* group = ctx->group; + if (ctx->timer) { uasync_cancel_timeout(group->instance->ua, ctx->timer); ctx->timer = NULL; } + if (ctx->wake) { uasync_call_soon_cancel(group->instance->ua, ctx->wake); ctx->wake = NULL; } + size_t remaining = recovery_remaining(ctx); + if (!remaining) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery complete: group=%016llx targets=%zu attempts=%u", + (unsigned long long)group->group_id, ctx->count, ctx->attempts); + recovery_finish(ctx); return; + } + if (ctx->request) { + enum topo_peer_phase phase = topo_group_peer_phase(ctx->request); + uint64_t now = get_time_tb(); + if (phase != ctx->phase) { + ctx->phase = phase; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery phase: group=%016llx peer=%016llx phase=%d missing=%zu", + (unsigned long long)group->group_id, (unsigned long long)ctx->current_node, phase, remaining); } - struct TOPO_RECOVERY_NODE* best = topo_recovery_find_best(ctx); - uint64_t node_id = best->node_id; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx trying %016llx rtt=%u %s", - (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id, best->min_rtt, - best->is_next_hop ? "[next_hop]" : ""); - ctx->current_node_id = node_id; - if (node_conn_direct_open(ctx->instance, node_id, topo_recovery_callback, ctx, - &ctx->current_handle, NULL) == NCD_ERR) { - DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: node_conn_direct_open %016llx failed", (unsigned long long)node_id); - topo_recovery_remove_candidate(ctx, (size_t)(best - ctx->nodes)); - ctx->current_node_id = 0; - continue; + uint64_t deadline = phase == TOPO_PEER_CONNECTING ? ctx->started_at + TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10ULL : + topo_group_peer_progress(ctx->request) + TOPO_RECOVERY_SYNC_TIMEOUT_MS * 10ULL; + if ((phase == TOPO_PEER_CONNECTING || phase == TOPO_PEER_SYNCING) && now < deadline) { + uint64_t delay = deadline - now; + ctx->timer = uasync_set_timeout(group->instance->ua, delay > INT_MAX ? INT_MAX : (int)delay, + ctx, recovery_timeout, "group_recovery"); + if (ctx->timer) return; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery timer allocation failed group=%016llx", (unsigned long long)group->group_id); } - if (!ctx->instance || !ctx->instance->ua) { - DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: instance/ua gone, aborting try_next for %016llx", - (unsigned long long)ctx->next_hop_id); - return; + if (phase == TOPO_PEER_CONNECTING || phase == TOPO_PEER_SYNCING) + DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery stalled: group=%016llx peer=%016llx phase=%d missing=%zu", + (unsigned long long)group->group_id, (unsigned long long)ctx->current_node, phase, remaining); + else DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery attempt ended: group=%016llx peer=%016llx phase=%d missing=%zu", + (unsigned long long)group->group_id, (unsigned long long)ctx->current_node, phase, remaining); + topo_group_peer_close(ctx->request); ctx->request = NULL; ctx->current_node = 0; + } + struct recovery_node* best; + while ((best = recovery_best(ctx)) != NULL) { + best->tried = 1; ctx->attempts++; + ctx->current_node = best->node_id; ctx->started_at = get_time_tb(); ctx->phase = TOPO_PEER_CONNECTING; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery attempt: group=%016llx peer=%016llx priority=%u rtt=%u missing=%zu attempt=%u", + (unsigned long long)group->group_id, (unsigned long long)best->node_id, best->priority, best->rtt, + remaining, ctx->attempts); + if (topo_group_peer_open(group, best->node_id, &ctx->request) == 0) { + /* Начальный UP готового NCD и READY существующей сессии проверяем асинхронно. */ + ctx->timer = uasync_set_timeout(group->instance->ua, 1, ctx, recovery_timeout, "group_recovery_start"); + if (ctx->timer) return; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery start timer allocation failed"); + topo_group_peer_close(ctx->request); ctx->request = NULL; } - ctx->connect_timer = uasync_set_timeout(ctx->instance->ua, TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10, - ctx, topo_recovery_timeout_cb, "recovery_timeout"); - return; } + DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery exhausted: group=%016llx missing=%zu attempts=%u", + (unsigned long long)group->group_id, remaining, ctx->attempts); + for (size_t i = 0; i < ctx->count; i++) + if (!ctx->nodes[i].restored) DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery unresolved: group=%016llx node=%016llx", + (unsigned long long)group->group_id, (unsigned long long)ctx->nodes[i].node_id); + recovery_finish(ctx); } -void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t next_hop) { - if (!group || !nq) return; - struct TOPO_NODE* ni = topo_node_registry_find(group->instance->topo_groups, nq->node_id); - if (!ni) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: missing candidate=%016llx", (unsigned long long)nq->node_id); return; } - uint64_t node_id = nq->node_id; - - struct TOPO_RECOVERY_CTX* ctx = topo_recovery_find_by_next_hop(group, next_hop); +void topo_recovery_add_node(struct TOPO_GROUP* group, uint64_t node_id, uint64_t next_hop, uint16_t rtt) { + if (!group || group->stopping) return; + if (!topo_node_registry_find(group->instance->topo_groups, node_id)) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery missing identity node=%016llx", (unsigned long long)node_id); return; + } + struct TOPO_RECOVERY_CTX* ctx = group->recovery; if (!ctx) { - ctx = topo_recovery_ctx_alloc(group, next_hop); - if (!ctx) return; - ctx->next = group->recovery_list; - group->recovery_list = ctx; + ctx = u_calloc(1, sizeof(*ctx)); + if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery context allocation failed"); return; } + ctx->group = group; group->recovery = ctx; } - - for (size_t i = 0; i < ctx->count; i++) if (ctx->nodes[i].node_id == node_id) return; - uint16_t rtt = topo_get_chain_rtt(nq); - - if (ctx->count >= ctx->capacity) { - size_t new_cap = ctx->capacity * 2; - struct TOPO_RECOVERY_NODE* new_nodes = u_realloc(ctx->nodes, new_cap * sizeof(struct TOPO_RECOVERY_NODE)); - if (!new_nodes) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery: candidate allocation failed"); return; } - ctx->nodes = new_nodes; ctx->capacity = new_cap; + for (size_t i = 0; i < ctx->count; i++) { + if (ctx->nodes[i].node_id != node_id) continue; + if (rtt < ctx->nodes[i].rtt) ctx->nodes[i].rtt = rtt; + if (node_id == next_hop) ctx->nodes[i].priority = 1; + return; + } + if (ctx->count == ctx->capacity) { + size_t capacity = ctx->capacity ? ctx->capacity * 2 : 16; + struct recovery_node* nodes = u_realloc(ctx->nodes, capacity * sizeof(*nodes)); + if (!nodes) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery candidate allocation failed"); return; } + ctx->nodes = nodes; ctx->capacity = capacity; } - topo_node_registry_ref(group->instance->topo_groups, node_id); - ctx->nodes[ctx->count].node_id = node_id; - ctx->nodes[ctx->count].min_rtt = rtt; - ctx->nodes[ctx->count].is_next_hop = (uint8_t)(node_id == next_hop ? 1 : 0); - ctx->count++; - - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "recovery: next=%016llx add node %016llx rtt=%u is_nh=%d count=%zu", - (unsigned long long)next_hop, (unsigned long long)node_id, - rtt, node_id == next_hop ? 1 : 0, ctx->count); + ctx->nodes[ctx->count++] = (struct recovery_node){ .node_id = node_id, .rtt = rtt, .priority = node_id == next_hop }; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery target: group=%016llx node=%016llx next=%016llx rtt=%u targets=%zu", + (unsigned long long)group->group_id, (unsigned long long)node_id, (unsigned long long)next_hop, rtt, ctx->count); } 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) { - topo_recovery_ctx_remove(group, ctx); - topo_recovery_ctx_free(ctx); - } else { - ctx->started = 1; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: start next=%016llx count=%zu", - (unsigned long long)ctx->next_hop_id, ctx->count); - 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"); - } + if (!group || !group->recovery || group->stopping) return; + group->recovery->started = 1; + topo_recovery_changed(group); } -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 (size_t i = 0; i < ctx->count && !found; i++) - if (ctx->nodes[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 && ctx->current_handle) - node_conn_direct_close(ctx->current_handle); - 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_changed(struct TOPO_GROUP* group) { + struct TOPO_RECOVERY_CTX* ctx = group ? group->recovery : NULL; + if (!ctx || !ctx->started || group->stopping || ctx->wake) return; + ctx->wake = uasync_call_soon(group->instance->ua, ctx, recovery_wake); + if (!ctx->wake) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery wake allocation failed group=%016llx", (unsigned long long)group->group_id); } 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 && ctx->current_handle) - node_conn_direct_close(ctx->current_handle); - 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; + if (!group || !group->recovery) return; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery cancelled group=%016llx", (unsigned long long)group->group_id); + recovery_finish(group->recovery); } diff --git a/src/routing_layer/topo_recovery.h b/src/routing_layer/topo_recovery.h index 7a503e6b..78e370d5 100644 --- a/src/routing_layer/topo_recovery.h +++ b/src/routing_layer/topo_recovery.h @@ -1,83 +1,42 @@ -/** - * @file topo_recovery.h - * @brief Восстановление каскадно отвалившихся узлов после разрыва ETCP-соединения. - * - * Когда рвётся соединение, узлы достижимые только через него становятся недоступны. - * Модуль собирает их перед удалением из BGP-таблицы и пробует переподключиться - * напрямую (ncd): - * 1. add_node — из hop_list вычисляет next_hop относительно failed_peer, - * группирует узлы по next_hop в отдельные recovery-контексты - * 2. start — для каждого контекста запускает асинхронный перебор: - * сперва пробует next_hop-узел (is_next_hop=1), затем остальные по мин. RTT - * 3. Таймаут 2с на каждую попытку node_conn_direct_open - * - * Группировка: <мы> -> -> - * Каждый next_hop со своим subtree — отдельный recovery-контекст. - * Восстановление next_hop автоматически оживляет его subtree через BGP. - * - * Отмена: topo_group_new_conn (узел появился в сети) сканирует все active recovery, - * при совпадении отменяет контекст целиком. Самоуничтожение при исчерпании списков. - */ #ifndef TOPO_RECOVERY_H #define TOPO_RECOVERY_H +#include + #ifdef __cplusplus extern "C" { #endif -#include -#include - struct TOPO_GROUP; -struct TOPO_GROUP_NODE; -struct NODE_CONN_DIRECT; +struct TOPO_RECOVERY_CTX; #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; /* кандидаты на восстановление (прямые) */ - size_t count; /* текущее количество */ - size_t capacity; /* выделенная ёмкость */ - uint64_t current_node_id; /* node_id в текущей попытке, 0=нет активной */ - struct NODE_CONN_DIRECT* current_handle; /* handle текущей попытки node_conn_direct_open */ - 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_GROUP_NODE* nq, uint64_t next_hop); - -/** - * Запускает перебор узлов для всех pending-контекстов. - * Вызывается после цикла в topo_group_remove_conn если есть каскадные узлы. - */ +#define TOPO_RECOVERY_SYNC_TIMEOUT_MS 5000 + +/* Один recovery-контекст на группу, одна текущая попытка присоединения. + * add_node сохраняет цель и кандидатуру до удаления registry ref; + * start вызывается после удаления всех путей через потерянного пира. + * Приоритет: бывшие downstream next-hop, затем минимальный RTT, затем node_id. + * Кандидат проверяется не более одного раза за цикл. Транспорт принадлежит + * групповой сессии: recovery владеет только отменяемым TOPO_PEER_REQUEST. + * + * Успех — все потерянные узлы снова имеют живой путь через READY-пира этой группы. + * UP и даже READY без нужных маршрутов не завершают recovery. Частичный результат + * сохраняет полезное присоединение и продолжает оставшиеся цели. Маршрут через + * другого READY-пира также закрывает цель. При исчерпании кандидатов оставшиеся + * цели логируются, контекст завершается без скрытого повторного цикла. + * + * CONNECTING ограничен отдельным таймаутом; SYNCING — временем без прогресса. + * NODEINFO/смена состояния лишь планируют проверку через call_soon: текущий + * BGP callback завершается до изменения попытки или освобождения контекста. + * Остановка группы отменяет recovery до освобождения сессий и таблицы узлов. */ +void topo_recovery_add_node(struct TOPO_GROUP* group, uint64_t node_id, uint64_t next_hop, uint16_t rtt); 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_changed(struct TOPO_GROUP* group); void topo_recovery_cancel_all(struct TOPO_GROUP* group); #ifdef __cplusplus } #endif + #endif diff --git a/tests/Makefile.am b/tests/Makefile.am index 00860404..438c042f 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -66,6 +66,7 @@ check_PROGRAMS = \ test_node_conn_direct \ test_group_ownership \ test_group_exchange \ + test_group_recovery \ test_node_snapshot \ test_ncd_config \ test_db_sync \ @@ -399,6 +400,10 @@ test_group_exchange_SOURCES = test_group_exchange.c test_group_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_group_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_group_recovery_SOURCES = test_group_recovery.c +test_group_recovery_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_group_recovery_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_node_snapshot_SOURCES = test_node_snapshot.c test_node_snapshot_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_node_snapshot_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_lifecycle.c b/tests/test_etcp_lifecycle.c index 33b68dc3..69c6cf82 100644 --- a/tests/test_etcp_lifecycle.c +++ b/tests/test_etcp_lifecycle.c @@ -111,15 +111,14 @@ static void recovery_cases(struct UTUN_INSTANCE* inst) { for (int scenario=0;scenario<3;scenario++) { struct TOPO_NODE* ni = u_calloc(1,sizeof(*ni)); CHECK(ni); ni->node_id = node.node_id; CHECK(topo_node_registry_store(&groups,ni) == ni && ni->group_ref_count == 1); - topo_recovery_add_node(&group,&node,44); topo_recovery_add_node(&group,&node,44); - CHECK(ni->group_ref_count == 2 && group.recovery_list->count == 1); - if (scenario == 1) { topo_recovery_add_node(&group,&node,45); CHECK(ni->group_ref_count == 3); } + topo_recovery_add_node(&group,node.node_id,44,100); topo_recovery_add_node(&group,node.node_id,44,100); + CHECK(ni->group_ref_count == 2 && group.recovery); + if (scenario == 1) { topo_recovery_add_node(&group,node.node_id,45,50); CHECK(ni->group_ref_count == 2); } topo_node_registry_unref(&groups,node.node_id); CHECK(topo_node_registry_find(&groups,node.node_id) == ni); - if (scenario == 0) topo_recovery_cancel_for_node(&group,node.node_id); - else if (scenario == 1) topo_recovery_cancel_all(&group); + if (scenario < 2) topo_recovery_cancel_all(&group); else { topo_recovery_start(&group); topo_recovery_cancel_all(&group); } uasync_poll(inst->ua,0); - CHECK(!group.recovery_list && !groups.node_registry->count && !inst->ncd_registry && !inst->connections->count); + CHECK(!group.recovery && !groups.node_registry->count && !inst->ncd_registry && !inst->connections->count); printf("PASS recovery ownership scenario %d\n",scenario+1); } queue_free(groups.node_registry); inst->topo_groups = NULL; diff --git a/tests/test_group_ownership.c b/tests/test_group_ownership.c index 7da0e0f0..5d41cd21 100644 --- a/tests/test_group_ownership.c +++ b/tests/test_group_ownership.c @@ -23,7 +23,7 @@ static void check_deliberate_leave(struct TOPO_GROUP* group, struct ETCP_CONN* c assert(queue_data_put_with_index(group->nodes, entry) == 0); topo_group_remove_conn(group, conn, reason); assert(!topo_node_find_by_id(group, id)); - assert(!group->recovery_list && !topo_node_registry_find(group->instance->topo_groups, id)); + assert(!group->recovery && !topo_node_registry_find(group->instance->topo_groups, id)); assert(!conn->close_requested); } @@ -51,10 +51,22 @@ int main(void) { /* Создание группы при уже существующем транспорте читает payload записи conn. */ struct TOPO_GROUP* c = topo_groups_create_group(inst->topo_groups, 44, TOPO_GROUP_TYPE_UTUN, NULL); assert(c); assert(queue_entry_count(c->senders_list) == 0); + struct TOPO_PEER_REQUEST *first = NULL, *second = NULL; + assert(topo_group_peer_open(c, node.node_id, &first) == 0); + assert(topo_group_peer_open(c, node.node_id, &second) == 0); + assert(topo_group_peer_phase(first) == TOPO_PEER_CONNECTING && topo_group_peer_phase(second) == TOPO_PEER_CONNECTING); + assert(queue_entry_count(c->senders_list) == 1); + topo_group_peer_close(first); + assert(queue_entry_count(c->senders_list) == 1 && !conn->close_requested); + topo_group_peer_close(second); + assert(queue_entry_count(c->senders_list) == 0 && !conn->close_requested); check_deliberate_leave(c, conn, TOPO_REMOVE_LOCAL_LEAVE); check_deliberate_leave(c, conn, TOPO_REMOVE_REMOTE_LEAVE); check_deliberate_leave(c, conn, TOPO_REMOVE_MEMBER_INVALID); + assert(topo_group_peer_open(c, node.node_id, &first) == 0); topo_groups_remove_group(inst->topo_groups, 44); + assert(topo_group_peer_phase(first) == TOPO_PEER_FAILED); + topo_group_peer_close(first); /* Проверяем владение независимо от доставки пакетов и ответа удалённого узла. */ assert(topo_group_new_conn(a, conn) == 0); assert(topo_group_new_conn(a, conn) == 0 && queue_entry_count(a->senders_list) == 1); diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c new file mode 100644 index 00000000..adab6bf2 --- /dev/null +++ b/tests/test_group_recovery.c @@ -0,0 +1,196 @@ +#include +#include +#include "utun_instance.h" +#include "topo_group.h" +#include "topo_recovery.h" +#include "etcp.h" +#include "etcp_api.h" +#include "node_conn_direct.h" +#include "../lib/mem.h" +#include "../lib/debug_config.h" +#include "../lib/platform_compat.h" + +struct fixture { + struct UASYNC* ua; + struct UTUN_INSTANCE* inst; + struct TOPO_GROUP* group; + uint64_t ids[3]; + struct ETCP_CONN* conns[3]; +}; + +static struct TOPO_GROUP_CONN_ITEM* peer(struct fixture* f, int index) { + for (struct ll_entry* e = f->group->senders_list->head; e; e = e->next) { + struct TOPO_GROUP_CONN_ITEM* p = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (p->node_id == f->ids[index]) return p; + } + return NULL; +} + +static void poll_events(struct fixture* f) { for (int i = 0; i < 4; i++) uasync_poll(f->ua, 0); } + +static void create(struct fixture* f) { + memset(f, 0, sizeof(*f)); + f->ua = uasync_create(); assert(f->ua); + f->inst = utun_instance_create_from_str(f->ua, + "[global]\n" + "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" + "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" + "[server: udp]\naddr=127.0.0.1:0\ntype=public\n"); + assert(f->inst && utun_instance_init(f->inst) == 0); + f->group = topo_groups_create_group(f->inst->topo_groups, 42, TOPO_GROUP_TYPE_UTUN, NULL); assert(f->group); + for (int i = 0; i < 3; i++) { + struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK); + struct TOPO_NODE* node = u_calloc(1, sizeof(*node)); assert(node); + memcpy(node->public_key, keys.public_key, SC_PUBKEY_SIZE); + f->ids[i] = node->node_id = sc_derive_node_id_from_pubkey(node->public_key); + assert(topo_node_registry_store(f->inst->topo_groups, node)); + struct ETCP_CONN* conn = f->conns[i] = etcp_connection_create(f->inst, "recovery_fixture"); assert(conn); + conn->peer_node_id = node->node_id; + queue_set_callback(conn->send_input_q, NULL, NULL); + etcp_conn_ready(conn); + } +} + +static void destroy(struct fixture* f) { + for (int i = 0; i < 3; i++) topo_node_registry_unref(f->inst->topo_groups, f->ids[i]); + f->inst->running = 0; utun_instance_destroy(f->inst); + uasync_poll(f->ua, 0); uasync_destroy(f->ua, 0); +} + +static void up(struct fixture* f, int i) { + f->conns[i]->links_up = 1; + etcp_cbk_fire(f->conns[i], ETCP_CBK_EVENT_UP); + assert(peer(f, i) && peer(f, i)->conn == f->conns[i]); +} + +static void route(struct fixture* f, int target, int via) { + struct ll_entry* entry = queue_entry_new(sizeof(struct TOPO_GROUP_NODE) - sizeof(struct ll_entry)); assert(entry); + struct TOPO_GROUP_NODE* node = (struct TOPO_GROUP_NODE*)entry; + node->node_id = f->ids[target]; + topo_node_registry_ref(f->inst->topo_groups, node->node_id); + uint64_t hops[] = { node->node_id, f->ids[via] }; + assert(topo_group_add_path(node, f->conns[via], hops, target == via ? 1 : 2, 10) == 0); + assert(queue_data_put_with_index(f->group->nodes, entry) == 0); + topo_recovery_changed(f->group); +} + +static void table(struct fixture* f, int i, uint8_t subcmd, uint64_t id) { + struct TOPOMSG_TABLE_REQ msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = subcmd, .group_id = f->group->group_id, .exchange_id = id }; + struct ll_entry* entry = ll_alloc_lldgram(sizeof(msg)); assert(entry); + memcpy(entry->dgram, &msg, sizeof(msg)); entry->len = sizeof(msg); + f->inst->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](f->conns[i], entry); +} + +static void ready(struct fixture* f, int i) { + assert(peer(f, i) && peer(f, i)->exchange_id); + table(f, i, TOPO_SUBCMD_REQUEST_TABLE, 7); + table(f, i, TOPO_SUBCMD_TABLE_COMPLETE, peer(f, i)->exchange_id); + assert(topo_group_peer_ready(f->group, f->ids[i])); +} + +static void partial_and_external_routes(void) { + struct fixture f; create(&f); + topo_recovery_add_node(f.group, f.ids[0], f.ids[0], 100); + topo_recovery_add_node(f.group, f.ids[1], f.ids[0], 1); + struct TOPO_RECOVERY_CTX* context = f.group->recovery; + topo_recovery_add_node(f.group, f.ids[1], f.ids[2], 1); + assert(context == f.group->recovery); + assert(topo_node_registry_find(f.inst->topo_groups, f.ids[1])->group_ref_count == 2); + topo_recovery_start(f.group); poll_events(&f); + assert(peer(&f, 0) && !peer(&f, 1) && !peer(&f, 2)); /* next-hop beats smaller RTT */ + topo_recovery_add_node(f.group, f.ids[2], f.ids[2], 200); + topo_recovery_start(f.group); poll_events(&f); + assert(f.group->recovery == context && peer(&f, 0) && !peer(&f, 2)); /* merge a new loss into the active round */ + struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle; + up(&f, 0); route(&f, 0, 0); poll_events(&f); + assert(f.group->recovery == context && !peer(&f, 2)); /* UP and a route are insufficient */ + assert(peer(&f, 0)->handle == ownership); /* no handoff or second NCD handle on UP */ + ready(&f, 0); poll_events(&f); + assert(f.group->recovery == context && peer(&f, 2) && !peer(&f, 1)); + assert(peer(&f, 0)->handle == ownership && topo_group_peer_ready(f.group, f.ids[0])); + route(&f, 1, 0); poll_events(&f); /* another READY peer restores the remaining subtree */ + up(&f, 2); route(&f, 2, 2); + table(&f, 2, TOPO_SUBCMD_TABLE_COMPLETE, peer(&f, 0)->exchange_id); poll_events(&f); + assert(f.group->recovery && !topo_group_peer_ready(f.group, f.ids[2])); + ready(&f, 2); poll_events(&f); + assert(!f.group->recovery && !peer(&f, 1)); + assert(topo_group_peer_ready(f.group, f.ids[0]) && topo_group_peer_ready(f.group, f.ids[2])); + destroy(&f); +} + +static void stalled_and_exhausted(void) { + struct fixture f; create(&f); + topo_recovery_add_node(f.group, f.ids[0], f.ids[0], 10); + topo_recovery_add_node(f.group, f.ids[1], f.ids[0], 20); + topo_recovery_start(f.group); poll_events(&f); + up(&f, 0); + uint64_t old_progress = get_time_tb() - TOPO_RECOVERY_SYNC_TIMEOUT_MS * 10ULL - 1; + peer(&f, 0)->progress = old_progress; + /* Реальный шаг обмена обновляет срок, пустая повторная проверка — нет. */ + table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, 8); poll_events(&f); + assert(peer(&f, 0) && !peer(&f, 1) && peer(&f, 0)->progress > old_progress); + peer(&f, 0)->progress = old_progress; + topo_recovery_changed(f.group); poll_events(&f); + assert(!peer(&f, 0) && peer(&f, 1)); + up(&f, 1); ready(&f, 1); poll_events(&f); /* READY without routes cannot claim success */ + assert(!f.group->recovery && !topo_node_find_by_id(f.group, f.ids[0]) && !topo_node_find_by_id(f.group, f.ids[1])); + topo_recovery_changed(f.group); topo_recovery_start(f.group); poll_events(&f); + assert(!f.group->recovery && !peer(&f, 0)); /* exhausted round does not restart itself */ + destroy(&f); +} + +static void cancel_and_group_stop(void) { + struct fixture f; create(&f); + struct TOPO_PEER_REQUEST* other = NULL; + assert(topo_group_peer_open(f.group, f.ids[0], &other) == 0); + struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle; + topo_recovery_add_node(f.group, f.ids[0], f.ids[0], 10); + topo_recovery_start(f.group); poll_events(&f); + assert(peer(&f, 0)->handle == ownership); + topo_recovery_cancel_all(f.group); poll_events(&f); + assert(!f.group->recovery && peer(&f, 0)->handle == ownership); + assert(topo_group_peer_phase(other) == TOPO_PEER_CONNECTING); + topo_recovery_add_node(f.group, f.ids[1], f.ids[1], 10); + topo_recovery_start(f.group); poll_events(&f); + topo_groups_remove_group(f.inst->topo_groups, f.group->group_id); f.group = NULL; + assert(topo_group_peer_phase(other) == TOPO_PEER_FAILED); + topo_group_peer_close(other); poll_events(&f); + destroy(&f); +} + +static void restored_during_connect(void) { + struct fixture f; create(&f); + f.conns[1]->links_up = 1; + assert(topo_group_new_conn(f.group, f.conns[1]) == 0); + route(&f, 1, 1); ready(&f, 1); + topo_recovery_add_node(f.group, f.ids[0], f.ids[0], 10); + topo_recovery_start(f.group); poll_events(&f); + assert(peer(&f, 0) && !peer(&f, 0)->conn); + route(&f, 0, 1); poll_events(&f); + assert(!f.group->recovery && !peer(&f, 0) && topo_group_peer_ready(f.group, f.ids[1])); + destroy(&f); +} + +static void connect_deadline(void) { + struct fixture f; create(&f); + f.inst->etcp_connect_timeout_tb = 100000; /* recovery must time out before NCD's own 10-second timer */ + topo_recovery_add_node(f.group, f.ids[0], f.ids[0], 10); + topo_recovery_start(f.group); poll_events(&f); + assert(peer(&f, 0)); + uint64_t deadline = get_time_tb() + (TOPO_RECOVERY_CONNECT_TIMEOUT_MS + 1000) * 10ULL; + while (f.group->recovery && get_time_tb() < deadline) uasync_poll(f.ua, 100); + assert(!f.group->recovery && !peer(&f, 0)); + destroy(&f); +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO); + utun_instance_set_tun_init_enabled(0); + partial_and_external_routes(); + stalled_and_exhausted(); + cancel_and_group_stop(); + restored_during_connect(); + connect_deadline(); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery tests: partial routes, readiness, progress, exhaustion and shared cancellation passed"); + return 0; +}