From 39ce249361dd072a23e01204434cc43dcd2826cb Mon Sep 17 00:00:00 2001 From: evgeny Date: Mon, 28 Sep 2026 14:58:13 +0300 Subject: [PATCH] Prune unconfirmed BGP paths and restart stalled exchanges --- src/routing_layer/topo_group.c | 77 ++++++++++++++++++++++++++++++++-- src/routing_layer/topo_group.h | 9 +++- tests/test_bgp_paths.c | 5 +++ tests/test_group_exchange.c | 23 +++++++++- tests/test_group_recovery.c | 2 + 5 files changed, 110 insertions(+), 6 deletions(-) diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 3dc9de81..c13944da 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -51,6 +51,23 @@ static int topo_tx_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t node_id); static int topo_tx_table(struct TOPO_GROUP_CONN_ITEM* peer); static int topo_control_wait(struct ETCP_CONN* conn, const void* data, size_t size); static void topo_group_routes_changed(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, int event); +static void topo_group_prune_paths(struct TOPO_GROUP_CONN_ITEM* peer, int all); +static void topo_exchange_watch(void* arg); +static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer); + +static void topo_exchange_cancel(struct TOPO_GROUP_CONN_ITEM* peer) { + if (peer->exchange_timer) uasync_cancel_timeout(peer->group->instance->ua, peer->exchange_timer); + peer->exchange_timer = NULL; +} + +static int topo_exchange_arm(struct TOPO_GROUP_CONN_ITEM* peer) { + if (!peer->exchange_timer) + peer->exchange_timer = uasync_set_timeout(peer->group->instance->ua, 1000, peer, topo_exchange_watch, "bgp_exchange"); + if (peer->exchange_timer) return 0; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot schedule exchange watchdog group=%016llx peer=%016llx", + (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id); + return -1; +} struct topo_advert { struct ll_entry ll; @@ -76,10 +93,12 @@ static struct TOPO_GROUP_CONN_ITEM* topo_group_peer(const struct TOPO_GROUP* gro 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 && peer->conn->links_up && !peer->conn->close_requested && peer->accepted && peer->table_received && peer->table_sent; + return peer && peer->conn && peer->conn->links_up && !peer->conn->close_requested && !peer->tx_failed && !peer->transport_down && + peer->accepted && peer->table_received && peer->table_sent; } static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { + topo_exchange_cancel(peer); topo_tx_clear(peer); topo_adverts_clear(peer); while (peer->requests) { @@ -190,6 +209,7 @@ static void topo_group_peer_progressed(struct TOPO_GROUP* group, uint64_t node_i 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)) { + topo_exchange_cancel(peer); peer->retained = 1; topo_group_connect_on_up(group, peer->conn); } topo_group_peer_progressed(group, peer->node_id); @@ -240,7 +260,7 @@ static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { topo_group_peer_progressed(peer->group, peer->node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "group session started: group=%016llx peer=%016llx local=%016llx", (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)epoch); - return 0; + return topo_exchange_arm(peer); } static void topo_group_send_table_request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { @@ -258,6 +278,22 @@ static void topo_group_send_join_group(struct TOPO_GROUP* group, struct ETCP_CON topo_send_control(conn, &msg, sizeof(msg)); } +static void topo_exchange_watch(void* arg) { + struct TOPO_GROUP_CONN_ITEM* peer = arg; + peer->exchange_timer = NULL; + if (!peer->conn || !peer->conn->links_up || peer->conn->close_requested || peer->transport_down || peer->group->stopping) return; + if (topo_group_peer_ready(peer->group, peer->node_id)) return; + uint64_t idle = get_time_tb() - peer->progress; + if (idle >= TOPO_GROUP_SYNC_TIMEOUT_MS * 10ULL) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "exchange stalled: group=%016llx peer=%016llx epoch=%016llx idle_ms=%llu sent=%u received=%u failed=%u", + (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, + (unsigned long long)peer->local_epoch, (unsigned long long)(idle / 10), + peer->table_sent, peer->table_received, peer->tx_failed); + topo_group_prune_paths(peer, 1); + if (topo_group_reset_exchange(peer) == 0) topo_group_send_join_group(peer->group, peer->conn); + } else topo_exchange_arm(peer); +} + static void topo_group_send_resync(struct ETCP_CONN* conn) { struct TOPOMSG_RESYNC msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_RESYNC }; topo_send_control(conn, &msg, sizeof(msg)); @@ -459,10 +495,15 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry* DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "BGP recv %s peer=%016llx group=%016llx local=%016llx remote=%016llx len=%u", group_subcmd_name(subcmd), (unsigned long long)peer->node_id, (unsigned long long)h->group_id, (unsigned long long)peer->local_epoch, (unsigned long long)peer->peer_epoch, entry->len); - if (subcmd == TOPO_SUBCMD_NODEINFO) { nodeinfo_dump_log(data, entry->len); topo_group_process_nodeinfo(group, from_conn, data, entry->len); } + if (peer->tx_failed) goto done; + if (subcmd == TOPO_SUBCMD_NODEINFO) { + nodeinfo_dump_log(data, entry->len); + if (topo_group_process_nodeinfo(group, from_conn, data, entry->len) < 0) topo_tx_fail(peer); + } else if (subcmd == TOPO_SUBCMD_WITHDRAW) topo_group_process_withdraw(group, from_conn, data, entry->len); else if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) topo_group_handle_request_table(group, from_conn); else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE && !peer->table_received) { + topo_group_prune_paths(peer, 0); peer->table_received = 1; topo_group_log_exchange(group, peer); } else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) { const struct TOPOMSG_ERR_GROUP_MISMATCH* err = (const struct TOPOMSG_ERR_GROUP_MISMATCH*)data; @@ -1011,6 +1052,7 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; if (item->conn == conn) { if (reason == TOPO_REMOVE_TRANSPORT_DOWN && item->requests && node_conn_direct_get_conn(item->handle) == conn) { + topo_exchange_cancel(item); topo_tx_clear(item); topo_adverts_clear(item); item->conn = NULL; item->transport_down = 1; item->retained = 0; item->accepted = 0; item->table_sent = 0; item->table_received = 0; topo_recovery_changed(group); topo_group_connect_changed(group); @@ -1113,6 +1155,29 @@ static void topo_group_routes_changed(struct TOPO_GROUP* group, struct TOPO_GROU topo_recovery_changed(group); } +/* TABLE_COMPLETE подтверждает только NODEINFO текущего обмена. Сбой обмена + * делает недостоверными все пути этого соседа, включая уже полученную часть. */ +static void topo_group_prune_paths(struct TOPO_GROUP_CONN_ITEM* peer, int all) { + struct TOPO_GROUP* group = peer->group; + for (struct ll_entry* e = group->nodes->head; e;) { + struct ll_entry* next = e->next; + struct TOPO_GROUP_NODE* node = (struct TOPO_GROUP_NODE*)e; + int remove = 0; + for (struct ll_entry* p = node->paths ? node->paths->head : NULL; p; p = p->next) { + struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)p; + if (path->conn == peer->conn && (all || path->exchange != peer->local_epoch)) remove = 1; + } + if (remove) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "prune peer route: group=%016llx node=%016llx peer=%016llx reason=%s", + (unsigned long long)group->group_id, (unsigned long long)node->node_id, + (unsigned long long)peer->node_id, all ? "exchange_timeout" : "not_in_snapshot"); + topo_group_remove_path(node, peer->conn); + topo_group_routes_changed(group, node, TOPO_NODE_EVENT_UPDATE); + } + e = next; + } +} + // ===== NODEINFO process ===== /* Шлёт пиру ошибку несоответствия типа группы (ERR_GROUP_MISMATCH). */ @@ -1277,7 +1342,10 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from u_free(new_hop_list); - if (topo_group_add_path(nodeinfo1, from, hop_list, extended_count, incoming_cumulative_rtt) < 0) return -1; + if (topo_group_add_path(nodeinfo1, from, hop_list, extended_count, incoming_cumulative_rtt) < 0) { + if (is_new_node) topo_nodeq_remove_node(group, nodeinfo1); + return -1; + } struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, from->peer_node_id); ((struct TOPO_NODEPATH*)nodeinfo1->paths->tail)->exchange = peer ? peer->local_epoch : 0; if (node_id != group->instance->node_id) { @@ -1440,6 +1508,7 @@ static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group sender failed: group=%016llx peer=%016llx epoch=%016llx", (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch); topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0; + topo_exchange_arm(peer); topo_recovery_changed(peer->group); topo_group_connect_changed(peer->group); } diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 1a1a9223..87de7773 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -118,7 +118,11 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g * не создаёт вторую таблицу. REINIT/leave/stop отменяют остаток и waiter; уже * переданные транспорту пакеты защищены поколениями на принимающей стороне. * - * Новое присоединение и REINIT создают новое поколение; REINIT сохраняет пути. + * Новое присоединение и REINIT создают новое поколение; REINIT временно сохраняет + * пути. Каждый NODEINFO помечает путь текущим local_epoch. На TABLE_COMPLETE + * удаляются неподтверждённые пути этого соседа; альтернативы остаются. + * Если обмен не продвигается TOPO_GROUP_SYNC_TIMEOUT_MS, все пути этого соседа + * удаляются и JOIN начинается заново. Ошибка применения/отправки запрещает READY. * Поколения защищают сообщения, а доступность существующего пути определяется * его соединением. Recovery дополнительно ожидает READY группового пира. * JOIN_READY чата означает добавление мембера и не заменяет этот обмен. @@ -181,6 +185,7 @@ struct TOPO_GROUP_CONN_ITEM { struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) struct TOPO_PEER_REQUEST* requests; // запросы инициаторов; каждый закрывает свой uint64_t progress; // время последнего прогресса обмена (timebase) + void* exchange_timer; // watchdog незавершённого обмена; отменяется при READY/DOWN/free uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch @@ -193,6 +198,8 @@ struct TOPO_GROUP_CONN_ITEM { uint8_t table_sent; // ответ на запрос пира целиком передан транспорту }; +#define TOPO_GROUP_SYNC_TIMEOUT_MS 5000 /* отсутствие прогресса: удалить пути пира и начать новый JOIN */ + struct ETCP_LINK; #define TOPO_NODE_REGISTRY_HASH_SIZE 256 diff --git a/tests/test_bgp_paths.c b/tests/test_bgp_paths.c index 8f64afd3..2b1a67a8 100644 --- a/tests/test_bgp_paths.c +++ b/tests/test_bgp_paths.c @@ -144,6 +144,11 @@ int main(void) { create(); unsigned triangle = edge(0, 1) | edge(1, 2) | edge(0, 2) | edge(2, 3); change(triangle); settle(); verify(); + /* One-sided REINIT with routes retained, then simultaneous REINIT at all peers. */ + etcp_fire_conn_status(conns[0][1], ETCP_CONN_STATUS_REINIT); settle(); verify(); + for (int a = 0; a < N; a++) for (int b = 0; b < N; b++) + if (graph & edge(a, b)) etcp_fire_conn_status(conns[a][b], ETCP_CONN_STATUS_REINIT); + settle(); verify(); change(triangle & ~edge(0, 2)); settle(); verify(); change(triangle); settle(); verify(); change(triangle & ~edge(1, 2) & ~edge(0, 2)); settle(); verify(); diff --git a/tests/test_group_exchange.c b/tests/test_group_exchange.c index 8b29175b..5fa188c8 100644 --- a/tests/test_group_exchange.c +++ b/tests/test_group_exchange.c @@ -8,6 +8,7 @@ #include "transport_layer/node_conn_direct.h" #include "../lib/u_async.h" #include "../lib/debug_config.h" +#include "../lib/platform_compat.h" static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) { assert(group->senders_list->head && !group->senders_list->head->next); @@ -153,8 +154,15 @@ int main(void) { assert(!peer(a)->table_sent); assert(peer(a)->tx_failed); conn->send_input_q->size_limit = saved_limit; - etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); + /* Watchdog must recover a failed sender without a transport status callback. */ + uint64_t failed_epoch = peer(a)->local_epoch; + peer(a)->progress = get_time_tb() - TOPO_GROUP_SYNC_TIMEOUT_MS * 10ULL - 1; + uint64_t deadline = get_time_tb() + 10000; + while (peer(a)->local_epoch == failed_epoch && get_time_tb() < deadline) uasync_poll(ua, 10); + assert(peer(a)->local_epoch != failed_epoch && !peer(a)->tx_failed); + assert(!topo_node_find_by_id(a, node.node_id)); /* incomplete snapshot expired */ accept_join(a, conn); + nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 1); complete(a, conn, peer(a)->local_epoch); request(a, conn); assert(topo_group_peer_ready(a, node.node_id)); @@ -169,6 +177,19 @@ int main(void) { assert(topo_node_find_by_id(a, node.node_id)); wd.h.dst_epoch = peer(a)->local_epoch; packet(conn, &wd, sizeof(wd)); assert(!topo_node_find_by_id(a, node.node_id)); + /* Same identity in a new exchange confirms a path; omission removes only that peer's path. */ + nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 2); + etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); accept_join(a, conn); + nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 2); + complete(a, conn, peer(a)->local_epoch); request(a, conn); + assert(topo_group_find_conn_for_node(a, node.node_id) == conn && !peer(a)->exchange_timer); + routed = topo_node_find_by_id(a, node.node_id); + assert(topo_group_add_path(routed, &relay, relay_hops, 2, 10) == 0); + etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); accept_join(a, conn); + complete(a, conn, peer(a)->local_epoch); request(a, conn); + assert(topo_group_find_conn_for_node(a, node.node_id) == &relay && queue_entry_count(routed->paths) == 1); + assert(topo_group_process_withdraw(a, &relay, (const uint8_t*)&alternate_wd, sizeof(alternate_wd)) == 0); + assert(!topo_node_find_by_id(a, node.node_id)); a_id = peer(a)->local_epoch; topo_group_remove_conn(a, conn, TOPO_REMOVE_LOCAL_LEAVE); complete(a, conn, a_id); diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c index 795daeea..cb528057 100644 --- a/tests/test_group_recovery.c +++ b/tests/test_group_recovery.c @@ -74,6 +74,8 @@ static void route(struct fixture* f, int target, int via) { 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); + /* This fixture injects the route that a NODEINFO of the current exchange would confirm. */ + if (peer(f, via)) ((struct TOPO_NODEPATH*)node->paths->tail)->exchange = peer(f, via)->local_epoch; assert(queue_data_put_with_index(f->group->nodes, entry) == 0); topo_recovery_changed(f->group); }