Browse Source

Reconcile BGP paths and advertise selected routes per neighbor

proxy
evgeny 3 days ago
parent
commit
b76e6553bc
  1. 286
      src/routing_layer/topo_group.c
  2. 40
      src/routing_layer/topo_group.h
  3. 46
      src/routing_layer/topo_node.c
  4. 3
      src/routing_layer/topo_node.h
  5. 5
      tests/Makefile.am
  6. 164
      tests/test_bgp_paths.c
  7. 4
      tests/test_group_exchange.c
  8. 3
      tests/test_group_recovery.c
  9. 1
      tests/test_node_snapshot.c

286
src/routing_layer/topo_group.c

@ -5,6 +5,7 @@
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <stdio.h> #include <stdio.h>
#include <openssl/sha.h>
#include "../lib/platform_compat.h" #include "../lib/platform_compat.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/mem.h" #include "../lib/mem.h"
@ -49,6 +50,21 @@ static int topo_tx_control(struct TOPO_GROUP_CONN_ITEM* peer, const void* data,
static int topo_tx_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t node_id); 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_tx_table(struct TOPO_GROUP_CONN_ITEM* peer);
static int topo_control_wait(struct ETCP_CONN* conn, const void* data, size_t size); 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);
struct topo_advert {
struct ll_entry ll;
uint64_t node_id;
uint8_t digest[SHA256_DIGEST_LENGTH];
};
/* Отправленные анонсы относятся только к текущему обмену с этим соседом. */
static void topo_adverts_clear(struct TOPO_GROUP_CONN_ITEM* peer) {
if (!peer->advertised) return;
struct ll_entry* e;
while ((e = queue_data_get(peer->advertised))) queue_entry_free(e);
queue_free(peer->advertised); peer->advertised = NULL;
}
static struct TOPO_GROUP_CONN_ITEM* topo_group_peer(const struct TOPO_GROUP* group, uint64_t peer_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) { for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) {
@ -65,6 +81,7 @@ int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id) {
static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) {
topo_tx_clear(peer); topo_tx_clear(peer);
topo_adverts_clear(peer);
while (peer->requests) { while (peer->requests) {
struct TOPO_PEER_REQUEST* request = peer->requests; struct TOPO_PEER_REQUEST* request = peer->requests;
peer->requests = request->next; request->peer = NULL; request->next = NULL; peer->requests = request->next; request->peer = NULL; request->next = NULL;
@ -211,6 +228,7 @@ static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t si
static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) {
topo_tx_clear(peer); peer->tx_failed = 0; peer->transport_down = 0; topo_tx_clear(peer); peer->tx_failed = 0; peer->transport_down = 0;
topo_adverts_clear(peer);
peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0; peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0;
uint64_t epoch; uint64_t epoch;
if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) { if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) {
@ -331,14 +349,12 @@ static const char* group_subcmd_name(uint8_t subcmd) {
// Broadcast / Withdraw // Broadcast / Withdraw
// ============================================================================ // ============================================================================
/* Рассылает WITHDRAW всем BGP-пирам (кроме exclude) об удалении узла. */ /* Пересчитать наш анонс назначения каждому соседу при освобождении его очереди. */
static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t node_id, uint64_t wd_source, struct ETCP_CONN* exclude) { static void topo_group_publish_route(struct TOPO_GROUP* group, uint64_t node_id) {
if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group is NULL"); return; } if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group is NULL"); return; }
for (struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; e; e = e->next) { for (struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; e; e = e->next) {
struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data; struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (!peer->conn || !peer->accepted || peer->conn == exclude) continue; if (peer->conn && peer->accepted && !peer->transport_down) topo_tx_node(peer, node_id);
struct TOPOMSG_WITHDRAW_PKT msg = { .h = topo_header(peer, TOPO_SUBCMD_WITHDRAW), .node_id = node_id, .wd_source = wd_source };
topo_send_control(peer->conn, &msg, sizeof(msg));
} }
} }
@ -945,7 +961,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en
(void*)conn, (unsigned long long)conn->peer_node_id, conn->state, conn->log_name, (void*)conn, (unsigned long long)conn->peer_node_id, conn->state, conn->log_name,
(unsigned long long)group->group_id, group->channel_id); (unsigned long long)group->group_id, group->channel_id);
struct ROUTE_TABLE* rt = conn->instance->rt; /* Больше не публикуем в закрываемую сессию, даже если общий транспорт ещё UP. */
pending->transport_down = 1;
int nodes_removed = 0; int nodes_removed = 0;
int cascaded = 0; int cascaded = 0;
struct ll_entry* node_entry = group->nodes ? group->nodes->head : NULL; struct ll_entry* node_entry = group->nodes ? group->nodes->head : NULL;
@ -969,34 +986,20 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en
} }
break; break;
} }
if (topo_group_remove_path(nq, conn) == 1) { int removed = topo_group_remove_path(nq, conn);
if (!removed) { node_entry = next; continue; }
if (!nq->paths) {
uint64_t key = nq->node_id; uint64_t key = nq->node_id;
if (reason == TOPO_REMOVE_TRANSPORT_DOWN && key != conn->peer_node_id) { if (reason == TOPO_REMOVE_TRANSPORT_DOWN && key != conn->peer_node_id) {
topo_recovery_add_node(group, key, next_hop, lost_rtt); 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);
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, key);
topo_nodeq_remove_node(group, nq);
nodes_removed++; nodes_removed++;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "node removed from group: node=%016llx grp=%016llx type=%u ch=%s (link_down)", DEBUG_INFO(DEBUG_CATEGORY_BGP, "node removed from group: node=%016llx grp=%016llx type=%u ch=%s (link_down)",
(unsigned long long)key, (unsigned long long)group->group_id, group->group_type, group->channel_id); (unsigned long long)key, (unsigned long long)group->group_id, group->group_type, group->channel_id);
// REINIT временно убирает путь, но end-to-end данные роутера остаются неподтверждёнными. // REINIT временно убирает путь, но end-to-end данные роутера остаются неподтверждёнными.
if (!conn->reinit_pending) etcp_router_conn_close_all_for_node(group->instance, group->group_id, key); if (!conn->reinit_pending) etcp_router_conn_close_all_for_node(group->instance, group->group_id, key);
uint64_t wd_src = (key == conn->peer_node_id) ? conn->instance->node_id : conn->peer_node_id;
topo_group_broadcast_withdraw(group, key, wd_src, NULL);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Removed node %016llx after link down", (unsigned long long)key);
{ struct topo_node_cbk_entry* c = group->node_cbks;
while (c) { c->fn(group, key, TOPO_NODE_EVENT_REMOVE, c->arg); c = c->next; }
}
} else {
nq->conn_up &= ~NCONN_DIRECT;
{ int has_direct = 0; struct ll_entry* pe2 = nq->paths ? nq->paths->head : NULL;
while (pe2) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)pe2; if (p->conn->peer_node_id == nq->node_id) { has_direct = 1; break; } pe2 = pe2->next; }
if (!has_direct) nq->conn_presence &= ~NCONN_DIRECT; }
topo_fire_nodeinfo_cbk(conn->instance, group, nq);
} }
topo_group_routes_changed(group, nq, TOPO_NODE_EVENT_UPDATE);
} }
node_entry = next; node_entry = next;
} }
@ -1008,7 +1011,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; struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data;
if (item->conn == conn) { if (item->conn == conn) {
if (reason == TOPO_REMOVE_TRANSPORT_DOWN && item->requests && node_conn_direct_get_conn(item->handle) == conn) { if (reason == TOPO_REMOVE_TRANSPORT_DOWN && item->requests && node_conn_direct_get_conn(item->handle) == conn) {
topo_tx_clear(item); item->conn = NULL; item->transport_down = 1; item->retained = 0; 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; item->accepted = 0; item->table_sent = 0; item->table_received = 0;
topo_recovery_changed(group); topo_group_connect_changed(group); topo_recovery_changed(group); topo_group_connect_changed(group);
} else { } else {
@ -1038,57 +1041,76 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en
/* Ищет лучшее ETCP-соединение до узла (мин. hop_count среди живых путей). */ /* Ищет лучшее ETCP-соединение до узла (мин. hop_count среди живых путей). */
struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id) { struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id) {
if (!group) return NULL; if (!group) return NULL;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id); struct TOPO_NODEPATH* path = topo_node_best_path(topo_node_find_by_id(group, node_id));
if (!nq || !nq->paths || !nq->paths->head) return NULL; return path ? path->conn : NULL;
struct ll_entry* e = nq->paths->head;
struct TOPO_NODEPATH* best_live = NULL; uint8_t best_live_hops = 255;
struct TOPO_NODEPATH* best_any = NULL; uint8_t best_any_hops = 255;
while (e) {
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)e;
if (path->conn) { if (path->conn->links_up && path->hop_count < best_live_hops) { best_live = path; best_live_hops = path->hop_count; } if (path->hop_count < best_any_hops) { best_any = path; best_any_hops = path->hop_count; } }
e = e->next;
}
if (!best_live && best_any) DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "find_conn_for_node %016llx: no live paths, fallback to best (hops=%d links_up=%d)", (unsigned long long)node_id, best_any_hops, best_any->conn->links_up);
return best_live ? best_live->conn : (best_any ? best_any->conn : NULL);
} }
/* Добавляет путь (conn + hop_list + rtt) в paths узла. */ /* Добавляет путь (conn + hop_list + rtt) в paths узла. */
int topo_group_add_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn, uint64_t* hop_list, uint8_t hop_count, uint16_t cumulative_rtt) { int topo_group_add_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn, uint64_t* hop_list, uint8_t hop_count, uint16_t cumulative_rtt) {
if (!nq || !conn || hop_count > MAX_HOPS || !hop_list) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "add_path: invalid args"); return -1; } if (!nq || !conn || !hop_count || hop_count > MAX_HOPS || !hop_list) {
if (!nq->paths) { nq->paths = queue_new(conn->instance->ua, 0, 0, 0, "node_paths"); if (!nq->paths) return -1; } DEBUG_ERROR(DEBUG_CATEGORY_BGP, "add_path: invalid args"); return -1;
}
size_t path_size = sizeof(struct TOPO_NODEPATH) - sizeof(struct ll_entry) + hop_count * 8; size_t path_size = sizeof(struct TOPO_NODEPATH) - sizeof(struct ll_entry) + hop_count * 8;
struct ll_entry* pe = queue_entry_new(path_size); if (!pe) return -1; struct ll_entry* pe = queue_entry_new(path_size);
if (!pe) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "path allocation failed"); return -1; }
if (!nq->paths) nq->paths = queue_new(conn->instance->ua, 0, 0, 0, "node_paths");
if (!nq->paths) { queue_entry_free(pe); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "paths queue allocation failed"); return -1; }
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe; struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe;
path->conn = conn; path->hop_count = hop_count; path->cumulative_rtt = cumulative_rtt; path->conn = conn; path->hop_count = hop_count; path->cumulative_rtt = cumulative_rtt;
memcpy((uint64_t*)((uint8_t*)path + sizeof(struct TOPO_NODEPATH)), hop_list, hop_count * 8); memcpy((uint64_t*)((uint8_t*)path + sizeof(struct TOPO_NODEPATH)), hop_list, hop_count * 8);
/* Все выделения завершены: прежний путь теперь можно заменить без потери при OOM. */
for (struct ll_entry* e = nq->paths->head; e;) {
struct ll_entry* next = e->next;
if (((struct TOPO_NODEPATH*)e)->conn == conn) { queue_remove_data(nq->paths, e); queue_entry_free(e); }
e = next;
}
queue_data_put(nq->paths, pe); queue_data_put(nq->paths, pe);
return 0; return 0;
} }
/* Удаляет пути, содержащие wd_source в hop_list; возвращает число удалённых путей. */ /* Сосед владеет только своим объявлением; возвращаем число удалённых путей. */
static int topo_group_remove_path_by_hop(struct TOPO_GROUP_NODE* nq, uint64_t wd_source) { int topo_group_remove_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn) {
if (!nq || !nq->paths) return 0; if (!nq || !conn || !nq->paths) return 0;
int removed = 0; int removed = 0;
struct ll_entry* e = nq->paths->head; for (struct ll_entry* e = nq->paths->head; e;) {
while (e) { struct ll_entry* next = e->next;
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)e; if (((struct TOPO_NODEPATH*)e)->conn == conn) { queue_remove_data(nq->paths, e); queue_entry_free(e); removed++; }
uint64_t* hop = (uint64_t*)((uint8_t*)path + sizeof(struct TOPO_NODEPATH)); e = next;
bool has_wd = false; for (uint8_t i = 0; i < path->hop_count; i++) if (hop[i] == wd_source) { has_wd = true; break; }
if (has_wd) { struct ll_entry* next = e->next; queue_remove_data(nq->paths, e); queue_entry_free(e); removed++; e = next; continue; }
e = e->next;
} }
if (removed > 0 && queue_entry_count(nq->paths) == 0) { queue_free(nq->paths); nq->paths = NULL; } if (removed > 0 && queue_entry_count(nq->paths) == 0) { queue_free(nq->paths); nq->paths = NULL; }
return removed; return removed;
} }
/* Удаляет путь через конкретный conn; 1 — если путей не осталось. */ /* Одна точка обновления доступности, потребителей и анонсов после изменения путей. */
int topo_group_remove_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn) { static void topo_group_routes_changed(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, int event) {
if (!nq || !conn || !nq->paths) return -1; uint64_t id = node->node_id;
struct ll_entry* e = nq->paths->head; struct TOPO_NODEPATH* best = topo_node_best_path(node);
while (e) { struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)e; if (path->conn == conn) { queue_remove_data(nq->paths, e); queue_entry_free(e); break; } e = e->next; } node->conn_presence &= ~(NCONN_BGP | NCONN_DIRECT);
int remaining = nq->paths ? queue_entry_count(nq->paths) : 0; node->conn_up &= ~(NCONN_BGP | NCONN_DIRECT);
if (remaining == 0) { if (nq->paths) { queue_free(nq->paths); nq->paths = NULL; } return 1; } if (node->paths) node->conn_presence |= NCONN_BGP;
return 0; if (best) node->conn_up |= NCONN_BGP;
for (struct ll_entry* e = node->paths ? node->paths->head : NULL; e; e = e->next) {
struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)e;
if (p->conn->peer_node_id != id) continue;
node->conn_presence |= NCONN_DIRECT;
if (p->conn->links_up && !p->conn->close_requested) node->conn_up |= NCONN_DIRECT;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "route changed: group=%016llx node=%016llx paths=%d via=%016llx hops=%u",
(unsigned long long)group->group_id, (unsigned long long)id, node->paths ? queue_entry_count(node->paths) : 0,
(unsigned long long)(best ? best->conn->peer_node_id : 0), best ? best->hop_count : 0);
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) {
if (best && node->subnets) route_insert(group->instance->rt, node);
else route_delete(group->instance->rt, node);
}
topo_fire_nodeinfo_cbk(group->instance, group, node);
if (!node->paths) {
event = TOPO_NODE_EVENT_REMOVE;
if (group->instance->control_srv) control_server_notify_node_removed(group->instance->control_srv, id);
topo_nodeq_remove_node(group, node);
} else if (group->instance->control_srv) control_server_notify_node_change(group->instance->control_srv, node);
for (struct topo_node_cbk_entry* c = group->node_cbks; c; c = c->next) c->fn(group, id, event, c->arg);
topo_group_publish_route(group, id);
topo_recovery_changed(group);
} }
// ===== NODEINFO process ===== // ===== NODEINFO process =====
@ -1178,6 +1200,22 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
u_free(new_hop_list); u_free(new_hop_list);
return -1; return -1;
} }
/* На проводе цепочка идёт от назначения до предыдущего узла, без отправителя.
* Нулевой hop_count допустим только для собственного анонса отправителя. */
int invalid_path = new_hop_count ? new_hop_list[0] != node_id : from->peer_node_id != node_id;
for (uint8_t i = 0; i < new_hop_count; i++) {
if (!new_hop_list[i] || new_hop_list[i] == group->instance->node_id || new_hop_list[i] == from->peer_node_id) invalid_path = 1;
for (uint8_t j = 0; j < i; j++) if (new_hop_list[j] == new_hop_list[i]) invalid_path = 1;
}
if (invalid_path) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "NODEINFO invalid path: group=%016llx node=%016llx from=%016llx hops=%u",
(unsigned long long)group->group_id, (unsigned long long)node_id,
(unsigned long long)from->peer_node_id, new_hop_count);
topo_node_destroy(group->instance->topo_groups, new_ni);
struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets };
topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected);
u_free(new_hop_list); return -1;
}
if (nodeinfo1 && new_ni->timestamp <= nodeinfo1->last_timestamp) { if (nodeinfo1 && new_ni->timestamp <= nodeinfo1->last_timestamp) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "NODEINFO unchanged/stale node=%016llx current=%llu incoming=%llu", DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "NODEINFO unchanged/stale node=%016llx current=%llu incoming=%llu",
(unsigned long long)node_id, (unsigned long long)nodeinfo1->last_timestamp, (unsigned long long)node_id, (unsigned long long)nodeinfo1->last_timestamp,
@ -1185,12 +1223,15 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
uint64_t hops[MAX_HOPS]; uint64_t hops[MAX_HOPS];
if (new_hop_count) memcpy(hops, new_hop_list, new_hop_count * sizeof(*hops)); if (new_hop_count) memcpy(hops, new_hop_list, new_hop_count * sizeof(*hops));
hops[new_hop_count] = from->peer_node_id; hops[new_hop_count] = from->peer_node_id;
topo_group_remove_path(nodeinfo1, from); int result = topo_group_add_path(nodeinfo1, from, hops, new_hop_count + 1, incoming_cumulative_rtt);
topo_group_add_path(nodeinfo1, from, hops, new_hop_count + 1, incoming_cumulative_rtt);
topo_node_destroy(group->instance->topo_groups, new_ni); topo_node_destroy(group->instance->topo_groups, new_ni);
struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets }; struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets };
topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected); topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected);
u_free(new_hop_list); u_free(new_hop_list);
if (result < 0) 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;
topo_group_routes_changed(group, nodeinfo1, TOPO_NODE_EVENT_UPDATE);
topo_group_peer_progressed(group, from->peer_node_id); topo_group_peer_progressed(group, from->peer_node_id);
return 0; return 0;
} }
@ -1236,13 +1277,9 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
u_free(new_hop_list); u_free(new_hop_list);
topo_group_remove_path_by_hop(nodeinfo1, from->peer_node_id); if (topo_group_add_path(nodeinfo1, from, hop_list, extended_count, incoming_cumulative_rtt) < 0) return -1;
topo_group_add_path(nodeinfo1, from, hop_list, extended_count, incoming_cumulative_rtt); 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;
nodeinfo1->conn_presence |= NCONN_BGP; nodeinfo1->conn_up |= NCONN_BGP;
if (from->peer_node_id == node_id) { nodeinfo1->conn_presence |= NCONN_DIRECT; nodeinfo1->conn_up |= NCONN_DIRECT; }
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) route_insert(group->instance->rt, nodeinfo1);
if (node_id != group->instance->node_id) { if (node_id != group->instance->node_id) {
sqlite3* sdb = group->instance->topo_sqlite_db; sqlite3* sdb = group->instance->topo_sqlite_db;
if (sdb) { if (sdb) {
@ -1260,34 +1297,11 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
topo_node_sqlite_nodeinfo_updated(sdb, node_id); topo_node_sqlite_nodeinfo_updated(sdb, node_id);
} }
} }
if (group->instance->control_srv) control_server_notify_node_change(group->instance->control_srv, nodeinfo1);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "node_updated_cb: node=%016llx has_cb=%d ch=%s type=%d", (unsigned long long)node_id, group->instance->topo_groups->node_updated_cb ? 1 : 0, group->channel_id, group->group_type); DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "node_updated_cb: node=%016llx has_cb=%d ch=%s type=%d", (unsigned long long)node_id, group->instance->topo_groups->node_updated_cb ? 1 : 0, group->channel_id, group->group_type);
if (group->instance->topo_groups->node_updated_cb) if (group->instance->topo_groups->node_updated_cb)
group->instance->topo_groups->node_updated_cb(group->instance, node_id, ni->public_key, ni->ed25519_public_key); group->instance->topo_groups->node_updated_cb(group->instance, node_id, ni->public_key, ni->ed25519_public_key);
{ /* fire BGP node event callbacks */ topo_group_routes_changed(group, nodeinfo1, is_new_node ? TOPO_NODE_EVENT_NEW : TOPO_NODE_EVENT_UPDATE);
int ev = is_new_node ? TOPO_NODE_EVENT_NEW : TOPO_NODE_EVENT_UPDATE;
struct topo_node_cbk_entry* c = group->node_cbks;
while (c) { c->fn(group, node_id, ev, c->arg); c = c->next; }
}
topo_fire_nodeinfo_cbk(group->instance, group, nodeinfo1);
int hop_count = extended_count;
struct ll_entry* se = group->senders_list ? group->senders_list->head : NULL;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item->conn) {
uint64_t id = item->conn->peer_node_id;
if (id != node_id) {
int found = 0; for (int i = 0; i < hop_count; i++) if (hop_list[i] == id) found = 1;
if (found == 0) {
topo_group_send_nodeinfo(group, nodeinfo1, item->conn);
}
}
}
se = se->next;
}
int socks_changed = 1; int socks_changed = 1;
if (!is_new_node) socks_changed = 1; if (!is_new_node) socks_changed = 1;
@ -1308,37 +1322,24 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
/* Обрабатывает WITHDRAW: удаляет узел/пути, чистит роутинг, форвардит withdraw. */ /* Обрабатывает WITHDRAW: удаляет узел/пути, чистит роутинг, форвардит withdraw. */
int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* sender, const uint8_t* data, size_t len) { int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* sender, const uint8_t* data, size_t len) {
if (!group || len < sizeof(struct TOPOMSG_WITHDRAW_PKT)) return -1; if (!group || !sender || len != sizeof(struct TOPOMSG_WITHDRAW_PKT)) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid WITHDRAW len=%zu", len); return -1;
}
struct TOPOMSG_WITHDRAW_PKT* wp = (struct TOPOMSG_WITHDRAW_PKT*)data; struct TOPOMSG_WITHDRAW_PKT* wp = (struct TOPOMSG_WITHDRAW_PKT*)data;
uint64_t node_id = wp->node_id, wd_source = wp->wd_source; uint64_t node_id = wp->node_id;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id); struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id);
if (!nq) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "node not found"); return 0; } if (!nq) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "WITHDRAW unknown node=%016llx", (unsigned long long)node_id); return 0; }
int removed = topo_group_remove_path_by_hop(nq, wd_source); int removed = topo_group_remove_path(nq, sender);
int remaining = nq->paths ? queue_entry_count(nq->paths) : 0; int remaining = nq->paths ? queue_entry_count(nq->paths) : 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "WITHDRAW group=%016llx node=%016llx source=%016llx removed=%d remaining=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "WITHDRAW group=%016llx node=%016llx source=%016llx removed=%d remaining=%d",
(unsigned long long)group->group_id, (unsigned long long)node_id, (unsigned long long)group->group_id, (unsigned long long)node_id,
(unsigned long long)wd_source, removed, remaining); (unsigned long long)sender->peer_node_id, removed, remaining);
if (removed > 0 && remaining == 0) { if (removed) topo_group_routes_changed(group, nq, TOPO_NODE_EVENT_UPDATE);
nq->conn_presence = 0; nq->conn_up = 0;
topo_fire_nodeinfo_cbk(group->instance, group, nq);
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance && group->instance->rt) route_delete(group->instance->rt, nq);
if (group->instance->control_srv) control_server_notify_node_removed(group->instance->control_srv, node_id);
topo_nodeq_remove_node(group, nq);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "node removed from group: node=%016llx grp=%016llx type=%u ch=%s (withdraw)",
(unsigned long long)node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id);
topo_group_broadcast_withdraw(group, node_id, wd_source, sender);
/* BGP не удаляет member-запись: он работает только со своей RAM-таблицей.
Удаление из peers_* — только битая подпись при верификации (member_sync). */
{ /* fire BGP REMOVE callback */
struct topo_node_cbk_entry* c = group->node_cbks;
while (c) { c->fn(group, node_id, TOPO_NODE_EVENT_REMOVE, c->arg); c = c->next; }
}
}
return 0; return 0;
} }
/* Сериализует и шлёт NODEINFO узла конкретному пиру. */ /* Сериализует актуальный анонс; кэш меняется только после успешной передачи ETCP. */
static int topo_emit_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt) { static int topo_emit_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt) {
if (!group || !node || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: invalid arguments"); return -1; } if (!group || !node || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: invalid arguments"); return -1; }
@ -1372,16 +1373,35 @@ static int topo_emit_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE*
u_free(p); return -1; u_free(p); return -1;
} }
if (!peer->advertised) peer->advertised = queue_new(group->instance->ua, BGP_NODES_HASH_SIZE, 0, 8, "topo_advertised");
if (!peer->advertised) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "advertised queue allocation failed"); u_free(p); return -1; }
uint8_t digest[SHA256_DIGEST_LENGTH];
SHA256(p + sizeof(h), ser_len, digest);
struct topo_advert* advert = (struct topo_advert*)queue_find_data_by_index(peer->advertised, &node->node_id);
if (advert && memcmp(advert->digest, digest, sizeof(digest)) == 0) { u_free(p); return 0; }
int fresh = advert == NULL;
if (fresh) advert = (struct topo_advert*)queue_entry_new(sizeof(*advert) - sizeof(struct ll_entry));
if (!advert) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "advert allocation failed"); u_free(p); return -1; }
struct ll_entry* e = queue_entry_new(0); struct ll_entry* e = queue_entry_new(0);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: entry allocation failed"); u_free(p); return -1; } if (!e) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: entry allocation failed");
if (fresh) queue_entry_free(&advert->ll);
u_free(p); return -1;
}
e->dgram = p; e->len = (size_t)ser_len + sizeof(h); e->dgram = p; e->len = (size_t)ser_len + sizeof(h);
if (etcp_send(conn, e) != 0) { if (etcp_send(conn, e) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s", DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s",
(unsigned long long)node->node_id, conn->log_name); (unsigned long long)node->node_id, conn->log_name);
if (fresh) queue_entry_free(&advert->ll);
u_free(p); queue_entry_free(e); return -1; u_free(p); queue_entry_free(e); return -1;
} }
return 0; advert->node_id = node->node_id; memcpy(advert->digest, digest, sizeof(digest));
if (fresh) queue_data_put_with_index(peer->advertised, &advert->ll);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "NODEINFO sent: group=%016llx node=%016llx to=%016llx",
(unsigned long long)group->group_id, (unsigned long long)node->node_id, (unsigned long long)peer->node_id);
return 1;
} }
/* NODEINFO сериализуется только при отправке: указатели на узлы не переживают callback. */ /* NODEINFO сериализуется только при отправке: указатели на узлы не переживают callback. */
@ -1491,8 +1511,17 @@ static int topo_tx_table(struct TOPO_GROUP_CONN_ITEM* peer) {
static int topo_tx_emit_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t id) { static int topo_tx_emit_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t id) {
struct TOPO_GROUP* group = peer->group; struct TOPO_GROUP* group = peer->group;
struct TOPO_GROUP_NODE* node = id == group->local_node->node_id ? group->local_node : topo_node_find_by_id(group, id); struct TOPO_GROUP_NODE* node = id == group->local_node->node_id ? group->local_node : topo_node_find_by_id(group, id);
if (!node || (node != group->local_node && !topo_group_should_send_to(node, peer->node_id))) return 0; if (!node || (node != group->local_node && !topo_group_should_send_to(node, peer->node_id))) {
return topo_emit_nodeinfo(group, node, peer->conn, node == group->local_node ? 0 : topo_get_chain_rtt(node)) < 0 ? -1 : 1; struct ll_entry* previous = peer->advertised ? queue_find_data_by_index(peer->advertised, &id) : NULL;
if (!previous) return 0;
struct TOPOMSG_WITHDRAW_PKT wd = { .h = topo_header(peer, TOPO_SUBCMD_WITHDRAW), .node_id = id };
if (topo_emit_control(peer->conn, &wd, sizeof(wd)) < 0) return -1;
queue_remove_data(peer->advertised, previous); queue_entry_free(previous);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "WITHDRAW sent: group=%016llx node=%016llx to=%016llx",
(unsigned long long)group->group_id, (unsigned long long)id, (unsigned long long)peer->node_id);
return 1;
}
return topo_emit_nodeinfo(group, node, peer->conn, node == group->local_node ? 0 : topo_get_chain_rtt(node));
} }
static void topo_tx_ready(struct ll_queue* queue, void* arg) { static void topo_tx_ready(struct ll_queue* queue, void* arg) {
@ -1615,17 +1644,13 @@ static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN*
/* Решает, надо ли слать узел пиру (не шлём, если пир уже есть в hop_list пути). */ /* Решает, надо ли слать узел пиру (не шлём, если пир уже есть в hop_list пути). */
static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t target_id) { static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t target_id) {
if (!nq || !nq->paths) return false; if (!nq) return false;
if (target_id == nq->node_id) return false; if (target_id == nq->node_id) return false;
struct ll_entry* e = nq->paths->head; struct TOPO_NODEPATH* path = topo_node_best_path(nq);
while (e) { if (!path || path->hop_count >= MAX_HOPS) return false;
struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)e; uint64_t* hop = (uint64_t*)((uint8_t*)path + sizeof(*path));
uint64_t* hop = (uint64_t*)((uint8_t*)path + sizeof(struct TOPO_NODEPATH)); for (uint8_t i = 0; i < path->hop_count; i++) if (hop[i] == target_id) return false;
bool has_id = false; for (uint8_t i = 0; i < path->hop_count; i++) if (hop[i] == target_id) { has_id = true; break; } return true;
if (!has_id) return true;
e = e->next;
}
return false;
} }
int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle) { int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle) {
@ -1800,10 +1825,3 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g
struct nodeinfo_cbk_entry* c = instance->nodeinfo_cbks; struct nodeinfo_cbk_entry* c = instance->nodeinfo_cbks;
while (c) { c->fn(group, node, c->arg); c = c->next; } while (c) { c->fn(group, node, c->arg); c = c->next; }
} }
/* Публичная рассылка WITHDRAW об удалении узла. */
void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id) {
if (!group) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "node=%016llx", (unsigned long long)node_id);
topo_group_broadcast_withdraw(group, node_id, group->instance->node_id, NULL);
}

40
src/routing_layer/topo_group.h

@ -8,6 +8,29 @@
* JOIN/ACCEPT согласует сессию, NODEINFO/WITHDRAW поддерживают её маршруты. * JOIN/ACCEPT согласует сессию, NODEINFO/WITHDRAW поддерживают её маршруты.
* Потерю путей обрабатывает topo_recovery, выбор CHAT-пиров — topo_group_connect. * Потерю путей обрабатывает topo_recovery, выбор CHAT-пиров — topo_group_connect.
* Все операции — в потоке uasync. Сервисы создают и удаляют свои группы явно. * Все операции — в потоке uasync. Сервисы создают и удаляют свои группы явно.
*
* Построение маршрутов:
* - Собственный NODEINFO имеет пустую цепочку. Получатель добавляет ID соседа
* в конец: цепочка хранится от назначения к нам, без нашего собственного ID.
* - У назначения не более одного пути от каждого соседа. NODEINFO атомарно
* заменяет его путь, WITHDRAW отзывает только его путь, DOWN удаляет все
* пути этого соединения. Чужой WITHDRAW не пересылается как глобальное удаление.
* - Timestamp версионирует подписанную запись узла, а не путь. Даже объявление
* со старой записью обновляет путь; сама запись при этом не откатывается.
* - После изменения выбираем живой путь с минимумом переходов; при равенстве
* побеждает меньший ID соседа. DOWN/closing-пути не используются как запасной результат.
* - Данные и NODEINFO используют один выбор пути. Выбранный путь объявляется
* всем соседям вне его цепочки. Если прежний анонс стал недопустим, шлём WITHDRAW.
* Цепочки с повторами или собственным ID на приёме отклоняются.
* - Каждая сессия хранит хеш последнего отправленного анонса назначения.
* Очередь хранит node_id; пакет строится из актуального состояния при отправке.
* Поэтому изменение пути распространяется без изменения timestamp и без эха дублей.
* - Альтернатива сохраняет доступность узла. После потери последнего пути
* удаляется только маршрут группы: членство и сохранённая запись узла остаются.
*
* После стабилизации связей обмен сходится в пределах связного графа группы
* и MAX_HOPS. READY относится к обмену с одним соседом. Временные петли данных
* при перестройке отсекает цепочка пройденных узлов пакета ETCP-router.
*/ */
#ifndef TOPO_GROUP_H #ifndef TOPO_GROUP_H
#define TOPO_GROUP_H #define TOPO_GROUP_H
@ -33,6 +56,7 @@ struct TOPO_RECOVERY_CTX;
struct TOPO_GROUP_CONNECT; struct TOPO_GROUP_CONNECT;
struct TOPO_PEER_REQUEST; struct TOPO_PEER_REQUEST;
struct topo_tx_item; struct topo_tx_item;
struct topo_advert;
struct broadcast_ctx; struct broadcast_ctx;
struct radio_ctx; struct radio_ctx;
@ -133,7 +157,7 @@ struct TOPOMSG_NODEINFO_PKT {
struct TOPOMSG_WITHDRAW_PKT { struct TOPOMSG_WITHDRAW_PKT {
struct TOPOMSG_HEADER h; struct TOPOMSG_HEADER h;
uint64_t node_id, wd_source; uint64_t node_id; /* отправитель отзывает только свой путь к этому узлу */
} __attribute__((packed)); } __attribute__((packed));
struct TOPOMSG_JOIN_GROUP { struct TOPOMSG_JOIN_GROUP {
@ -164,6 +188,7 @@ struct TOPO_GROUP_CONN_ITEM {
struct topo_tx_item *tx_head, *tx_tail; // FIFO команд, ещё не переданных транспорту struct topo_tx_item *tx_head, *tx_tail; // FIFO команд, ещё не переданных транспорту
struct queue_waiter_handle tx_waiter; // ожидание свободной send_input_q struct queue_waiter_handle tx_waiter; // ожидание свободной send_input_q
void* tx_wake; // отложенный запуск отправки void* tx_wake; // отложенный запуск отправки
struct ll_queue* advertised; // node_id -> хеш последнего успешно отправленного NODEINFO; владеет сессия
uint8_t table_received; // TABLE_COMPLETE текущей сессии принят uint8_t table_received; // TABLE_COMPLETE текущей сессии принят
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
}; };
@ -348,16 +373,11 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
*/ */
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn); int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn);
/**
* @brief Отправляет WITHDRAW для node_id (вызывает broadcast_withdraw).
*/
void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id);
/** /**
* @brief Поиск оптимального ETCP соединения для указанного node_id. * @brief Поиск оптимального ETCP соединения для указанного node_id.
* *
* Выбирает минимум hop_count среди живых путей. Если живых нет, возвращает * Выбирает минимум hop_count среди UP-путей, при равенстве — минимум ID соседа.
* лучший сохранённый путь: ненулевой результат сам по себе не гарантирует UP. * При отсутствии живого пути возвращает NULL.
* *
* @param group указатель на TOPO_GROUP * @param group указатель на TOPO_GROUP
* @param node_id целевой узел * @param node_id целевой узел
@ -366,7 +386,7 @@ void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id);
struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id); struct ETCP_CONN* topo_group_find_conn_for_node(struct TOPO_GROUP* group, uint64_t node_id);
/** /**
* @brief Добавляет путь (conn) в paths узла. * @brief Атомарно заменяет путь через conn; остальные соседи не затрагиваются.
* *
* @return 0 при успехе * @return 0 при успехе
*/ */
@ -375,7 +395,7 @@ int topo_group_add_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn, uint
/** /**
* @brief Удаляет conn из paths узла. * @brief Удаляет conn из paths узла.
* *
* @return 1 если путей не осталось (unreachable) * @return число удалённых путей; доступность проверяется отдельно
*/ */
int topo_group_remove_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn); int topo_group_remove_path(struct TOPO_GROUP_NODE* nq, struct ETCP_CONN* conn);

46
src/routing_layer/topo_node.c

@ -127,17 +127,24 @@ static int reality_sock_list_equal(struct TOPO_REALITY_SOCK* a, struct TOPO_REAL
return a == NULL && b == NULL; return a == NULL && b == NULL;
} }
/* извлекает hop_list из лучшего живого path в nq->paths (или любого, если живых нет). /* Один выбор для пересылки данных и NODEINFO. Порядок доставки не влияет на равные пути. */
память принадлежит TOPO_NODEPATH в paths, не освобождать */ struct TOPO_NODEPATH* topo_node_best_path(const struct TOPO_GROUP_NODE* nq) {
struct TOPO_NODEPATH* best = NULL;
for (struct ll_entry* e = nq && nq->paths ? nq->paths->head : NULL; e; e = e->next) {
struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)e;
if (!p->conn || !p->conn->links_up || p->conn->close_requested) continue;
if (!best || p->hop_count < best->hop_count || (p->hop_count == best->hop_count &&
p->conn->peer_node_id < best->conn->peer_node_id)) best = p;
}
return best;
}
/* Цепочка выбранного живого пути; память принадлежит paths, не освобождать. */
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt) { uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt) {
if (!nq || !nq->paths || !nq->paths->head) { *out_count = 0; if (out_rtt) *out_rtt = 0; return NULL; } struct TOPO_NODEPATH* best = topo_node_best_path(nq);
struct TOPO_NODEPATH* best = NULL; uint8_t min_hops = 255; *out_count = best ? best->hop_count : 0;
struct ll_entry* e = nq->paths->head; if (out_rtt) *out_rtt = best ? best->cumulative_rtt : 0;
while (e) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)e; if (p->conn && p->conn->links_up && p->hop_count < min_hops) { best = p; min_hops = p->hop_count; } e = e->next; } return best ? (uint64_t*)((uint8_t*)best + sizeof(*best)) : NULL;
if (!best) best = (struct TOPO_NODEPATH*)nq->paths->head;
*out_count = best->hop_count;
if (out_rtt) *out_rtt = best->cumulative_rtt;
return (uint64_t*)((uint8_t*)best + sizeof(struct TOPO_NODEPATH));
} }
/* Освобождает списки адресов (sock_meta + addrs + reality) узла в глобальном реестре. /* Освобождает списки адресов (sock_meta + addrs + reality) узла в глобальном реестре.
@ -708,21 +715,12 @@ void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uin
} }
} }
/* Суммарный RTT цепочки до узла: минимальный (rtt_last + cumulative_rtt) среди живых путей. */ /* RTT именно выбранного пути, с насыщением вместо переполнения. */
uint16_t topo_get_chain_rtt(struct TOPO_GROUP_NODE* nq) { uint16_t topo_get_chain_rtt(struct TOPO_GROUP_NODE* nq) {
if (!nq || !nq->paths || !nq->paths->head) { struct TOPO_NODEPATH* p = topo_node_best_path(nq);
return 0xFFFF; if (!p || !p->conn->rtt_last) return UINT16_MAX;
} uint32_t total = (uint32_t)p->conn->rtt_last + p->cumulative_rtt;
uint16_t best = 0xFFFF; return total > UINT16_MAX ? UINT16_MAX : (uint16_t)total;
for (struct ll_entry* e = nq->paths->head; e; e = e->next) {
struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)e;
if (!p->conn || !p->conn->links_up || !p->conn->rtt_last) {
continue;
}
uint16_t total = (uint16_t)((uint32_t)p->conn->rtt_last + p->cumulative_rtt);
if (total < best) best = total;
}
return best;
} }
// ===== dump / format ===== // ===== dump / format =====

3
src/routing_layer/topo_node.h

@ -191,6 +191,7 @@ struct TOPO_NODEPATH {
struct ETCP_CONN* conn; // заимствованный next-hop; транспорт удерживает сессия struct ETCP_CONN* conn; // заимствованный next-hop; транспорт удерживает сессия
uint8_t hop_count; // число node_id в массиве сразу после структуры uint8_t hop_count; // число node_id в массиве сразу после структуры
uint16_t cumulative_rtt; // RTT оставшейся цепочки, без локального линка (0.1 ms) uint16_t cumulative_rtt; // RTT оставшейся цепочки, без локального линка (0.1 ms)
uint64_t exchange; // local_epoch сессии, в которой путь подтверждён NODEINFO
}; };
/** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes (ll_queue). */ /** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes (ll_queue). */
@ -245,6 +246,8 @@ int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t
/* Заимствованный список hops: кратчайший живой путь, иначе первый сохранённый. /* Заимствованный список hops: кратчайший живой путь, иначе первый сохранённый.
* out_count обязателен; NULL означает отсутствие пути. */ * out_count обязателен; NULL означает отсутствие пути. */
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt); uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt);
/* Лучший живой путь: минимум переходов, затем ID соседа. Указатель заимствованный. */
struct TOPO_NODEPATH* topo_node_best_path(const struct TOPO_GROUP_NODE* nq);
/* Обновить локальный узел группы и его анонс. 0/-1. */ /* Обновить локальный узел группы и его анонс. 0/-1. */
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group); int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);

5
tests/Makefile.am

@ -55,6 +55,7 @@ check_PROGRAMS = \
test_socks_client \ test_socks_client \
test_bgp_route_exchange \ test_bgp_route_exchange \
test_bgp_triangle \ test_bgp_triangle \
test_bgp_paths \
test_broadcast \ test_broadcast \
test_conn_mgr \ test_conn_mgr \
test_conn_mgr_phases \ test_conn_mgr_phases \
@ -401,6 +402,10 @@ test_group_exchange_SOURCES = test_group_exchange.c
test_group_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_group_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_group_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_group_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_bgp_paths_SOURCES = test_bgp_paths.c
test_bgp_paths_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_bgp_paths_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_group_recovery_SOURCES = test_group_recovery.c test_group_recovery_SOURCES = test_group_recovery.c
test_group_recovery_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_group_recovery_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_group_recovery_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_group_recovery_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

164
tests/test_bgp_paths.c

@ -0,0 +1,164 @@
/* Управляемая сеть: настоящий BGP receiver/sender, подписи и очереди, без UDP.
* Доставка по каждому ребру сохраняет FIFO, порядок между рёбрами меняется. */
#include <assert.h>
#include <stdio.h>
#include <string.h>
#include "utun_instance.h"
#include "config_updater.h"
#include "topo_group.h"
#include "etcp.h"
#include "etcp_api.h"
#include "node_conn_direct.h"
#include "../lib/debug_config.h"
#include "test_utils.h"
#define N 4
#define GROUP 42
static struct UASYNC* ua;
static struct UTUN_INSTANCE* nodes[N];
static struct TOPO_GROUP* groups[N];
static struct ETCP_CONN* conns[N][N];
static struct NODE_CONN_DIRECT* owners[N][N];
static unsigned graph, delivered;
static uint32_t random_state = 1;
static unsigned next_random(void) {
random_state = random_state * 1664525U + 1013904223U;
return random_state;
}
static unsigned edge(int a, int b) {
if (a > b) { int t = a; a = b; b = t; }
unsigned bit = 1;
for (int i = 0; i < N; i++) for (int j = i + 1; j < N; j++, bit <<= 1)
if (i == a && j == b) return bit;
return 0;
}
static void release(struct ll_entry* e) { queue_dgram_free(e); queue_entry_free(e); }
/* Единственный ручной потребитель send_input_q; штатный normalizer отключён. */
static int step(void) {
uasync_poll(ua, 0);
int work = 0, offset = next_random() % (N * N);
for (int k = 0; k < N * N; k++) {
int a = ((k + offset) % (N * N)) / N, b = (k + offset) % N;
if (a == b) continue;
struct ll_queue* q = conns[a][b]->send_input_q;
assert(queue_entry_count(q) <= 1);
struct ll_entry* e = queue_data_get(q);
if (!e) continue;
work++;
if ((graph & edge(a, b)) && e->len && e->dgram[0] == ETCP_ID_TOPO_ENTRY) {
delivered++;
nodes[b]->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](conns[b][a], e);
} else release(e);
queue_resume_callback(q);
}
return work;
}
static void settle(void) {
unsigned before = delivered;
int idle = 0;
for (int i = 0; i < 4000 && idle < 16; i++) idle = step() ? 0 : idle + 1;
assert(idle == 16);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "model settled graph=%02x packets=%u", graph, delivered - before);
}
/* Смена графа без обработки очередей между событиями разрыва и восстановления. */
static void change(unsigned mask) {
unsigned old = graph; graph = mask;
for (int a = 0; a < N; a++) for (int b = a + 1; b < N; b++) {
if (!(old & edge(a, b)) || (mask & edge(a, b))) continue;
conns[a][b]->links_up = conns[b][a]->links_up = 0;
topo_group_remove_conn(groups[a], conns[a][b], TOPO_REMOVE_REMOTE_LEAVE);
topo_group_remove_conn(groups[b], conns[b][a], TOPO_REMOVE_REMOTE_LEAVE);
struct ll_entry* e;
while ((e = queue_data_get(conns[a][b]->send_input_q))) release(e);
while ((e = queue_data_get(conns[b][a]->send_input_q))) release(e);
}
for (int a = 0; a < N; a++) for (int b = a + 1; b < N; b++) {
if ((old & edge(a, b)) || !(mask & edge(a, b))) continue;
conns[a][b]->links_up = conns[b][a]->links_up = 1;
assert(topo_group_new_conn(groups[a], conns[a][b]) == 0);
assert(topo_group_new_conn(groups[b], conns[b][a]) == 0);
}
}
/* Независимый эталон — кратчайшие расстояния в графе реальных соединений. */
static void verify(void) {
int distance[N][N];
for (int a = 0; a < N; a++) for (int b = 0; b < N; b++)
distance[a][b] = a == b ? 0 : (graph & edge(a, b)) ? 1 : 100;
for (int k = 0; k < N; k++) for (int a = 0; a < N; a++) for (int b = 0; b < N; b++)
if (distance[a][k] + distance[k][b] < distance[a][b]) distance[a][b] = distance[a][k] + distance[k][b];
for (int a = 0; a < N; a++) for (int b = 0; b < N; b++) {
if (a == b) continue;
struct ETCP_CONN* conn = topo_group_find_conn_for_node(groups[a], nodes[b]->node_id);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "verify graph=%02x src=%d dst=%d expected=%d conn=%p", graph, a, b, distance[a][b], conn);
assert((conn != NULL) == (distance[a][b] < 100));
if (!conn) continue;
struct TOPO_GROUP_NODE* node = topo_node_find_by_id(groups[a], nodes[b]->node_id);
uint8_t count; uint64_t* hops = topo_node_best_hop_list(node, &count, NULL);
assert(count == distance[a][b]);
int previous = a;
for (int h = count - 1; h >= 0; h--) {
int current = 0;
while (current < N && nodes[current]->node_id != hops[h]) current++;
assert(current < N && (graph & edge(previous, current)));
previous = current;
}
assert(previous == b);
for (int h = 0; h < count; h++) for (int j = h + 1; j < count; j++) assert(hops[h] != hops[j]);
if (graph & edge(a, b)) assert(topo_group_peer_ready(groups[a], nodes[b]->node_id));
}
}
static void create(void) {
ua = uasync_create(); assert(ua);
for (int i = 0; i < N; i++) {
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
char pub[65], priv[65], config[512];
bytes_to_hex(keys.public_key, 32, pub, sizeof(pub)); bytes_to_hex(keys.private_key, 32, priv, sizeof(priv));
snprintf(config, sizeof(config), "[global]\nmy_public_key=%s\nmy_private_key=%s\n", pub, priv);
nodes[i] = utun_instance_create_from_str(ua, config); assert(nodes[i]);
assert(utun_core_start(nodes[i]) == 0);
groups[i] = topo_groups_create_group(nodes[i]->topo_groups, GROUP, TOPO_GROUP_TYPE_UTUN, NULL); assert(groups[i]);
}
for (int a = 0; a < N; a++) for (int b = 0; b < N; b++) {
if (a == b) continue;
struct ETCP_CONN* c = conns[a][b] = etcp_connection_create(nodes[a], "virtual_bgp"); assert(c);
c->peer_node_id = nodes[b]->node_id;
assert(sc_init_ctx(&c->crypto_ctx, &nodes[a]->my_keys) == SC_OK);
assert(sc_set_peer_public_key(&c->crypto_ctx, nodes[b]->my_keys.public_key, SC_PEER_PUBKEY_BIN) == SC_OK);
queue_set_callback(c->send_input_q, NULL, NULL);
etcp_conn_ready(c);
assert(node_conn_direct_open(nodes[a], c->peer_node_id, NULL, NULL, &owners[a][b], NULL) >= 0);
}
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); debug_set_category_level(DEBUG_CATEGORY_BGP, DEBUG_LEVEL_DEBUG);
utun_instance_set_tun_init_enabled(0);
create();
unsigned triangle = edge(0, 1) | edge(1, 2) | edge(0, 2) | edge(2, 3);
change(triangle); settle(); verify();
change(triangle & ~edge(0, 2)); settle(); verify();
change(triangle); settle(); verify();
change(triangle & ~edge(1, 2) & ~edge(0, 2)); settle(); verify();
change(triangle); settle(); verify();
/* Все графы, включая изолированные компоненты; затем churn без ожидания сходимости. */
for (unsigned mask = 0; mask < 64; mask++) { change(mask); settle(); verify(); }
for (int i = 0; i < 100; i++) {
change(next_random() >> 26);
for (unsigned n = next_random() % 5; n; n--) step();
}
settle(); verify();
change(0); settle(); verify();
for (int a = 0; a < N; a++) for (int b = 0; b < N; b++) if (a != b) node_conn_direct_close(owners[a][b]);
for (int a = 0; a < N; a++) utun_instance_destroy(nodes[a]);
uasync_poll(ua, 0); uasync_destroy(ua, 0);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP paths: triangle, partitions, all four-node graphs and churn passed packets=%u", delivered);
return 0;
}

4
tests/test_group_exchange.c

@ -131,7 +131,7 @@ int main(void) {
uint64_t relay_hops[] = { node.node_id, relay.peer_node_id }; uint64_t relay_hops[] = { node.node_id, relay.peer_node_id };
struct TOPO_GROUP_NODE* routed = topo_node_find_by_id(a, node.node_id); struct TOPO_GROUP_NODE* routed = topo_node_find_by_id(a, node.node_id);
assert(topo_group_add_path(routed, &relay, relay_hops, 2, 10) == 0); assert(topo_group_add_path(routed, &relay, relay_hops, 2, 10) == 0);
struct TOPOMSG_WITHDRAW_PKT alternate_wd = { .node_id = node.node_id, .wd_source = relay.peer_node_id }; struct TOPOMSG_WITHDRAW_PKT alternate_wd = { .node_id = node.node_id };
DEBUG_INFO(DEBUG_CATEGORY_BGP, "regression: withdraw one of two paths node=%016llx paths=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "regression: withdraw one of two paths node=%016llx paths=%d",
(unsigned long long)node.node_id, queue_entry_count(routed->paths)); (unsigned long long)node.node_id, queue_entry_count(routed->paths));
assert(topo_group_process_withdraw(a, &relay, (const uint8_t*)&alternate_wd, sizeof(alternate_wd)) == 0); assert(topo_group_process_withdraw(a, &relay, (const uint8_t*)&alternate_wd, sizeof(alternate_wd)) == 0);
@ -164,7 +164,7 @@ int main(void) {
assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 2); assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 2);
struct TOPOMSG_WITHDRAW_PKT wd = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_WITHDRAW, struct TOPOMSG_WITHDRAW_PKT wd = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_WITHDRAW,
.group_id = a->group_id, .src_epoch = a->group_id + 1000, .dst_epoch = a_id }, .group_id = a->group_id, .src_epoch = a->group_id + 1000, .dst_epoch = a_id },
.node_id = node.node_id, .wd_source = node.node_id }; .node_id = node.node_id };
packet(conn, &wd, sizeof(wd)); packet(conn, &wd, sizeof(wd));
assert(topo_node_find_by_id(a, node.node_id)); assert(topo_node_find_by_id(a, node.node_id));
wd.h.dst_epoch = peer(a)->local_epoch; packet(conn, &wd, sizeof(wd)); wd.h.dst_epoch = peer(a)->local_epoch; packet(conn, &wd, sizeof(wd));

3
tests/test_group_recovery.c

@ -245,7 +245,8 @@ static void sender_backpressure(void) {
if (e) { queue_dgram_free(e); queue_entry_free(e); } if (e) { queue_dgram_free(e); queue_entry_free(e); }
queue_resume_callback(f.conns[0]->send_input_q); queue_resume_callback(f.conns[0]->send_input_q);
} }
assert(complete == 1 && local == 2 && remote == 1 && peer(&f, 0)->table_sent); /* 50 одинаковых обновлений после снимка не повторяют уже переданный анонс. */
assert(complete == 1 && local == 1 && remote == 1 && peer(&f, 0)->table_sent);
/* Cancel a blocked old generation; only the new JOIN may be sent afterwards. */ /* Cancel a blocked old generation; only the new JOIN may be sent afterwards. */
etcp_fire_conn_status(f.conns[0], ETCP_CONN_STATUS_REINIT); accept_join(&f, 0); etcp_fire_conn_status(f.conns[0], ETCP_CONN_STATUS_REINIT); accept_join(&f, 0);
table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch); table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch);

1
tests/test_node_snapshot.c

@ -41,6 +41,7 @@ static void test_bgp(struct UTUN_INSTANCE* signer) {
struct ETCP_CONN from = { .instance = &receiver, .peer_node_id = 7 }; struct ETCP_CONN from = { .instance = &receiver, .peer_node_id = 7 };
group.nodes = queue_new(receiver.ua, 16, 0, 8, "snapshot_bgp"); group.nodes = queue_new(receiver.ua, 16, 0, 8, "snapshot_bgp");
struct TOPO_NODE* ni = record(signer, INT64_MAX - 7, "BGP", 1); struct TOPO_NODE* ni = record(signer, INT64_MAX - 7, "BGP", 1);
from.peer_node_id = ni->node_id; from.links_up = 1;
struct TOPO_GROUP_NODE nq = {0}; struct TOPO_GROUP_NODE nq = {0};
uint8_t packet[4096] = {0}; uint8_t packet[4096] = {0};
int len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + sizeof(struct TOPOMSG_HEADER), sizeof(packet) - sizeof(struct TOPOMSG_HEADER), 0); int len = topo_node_serialize(ni, &nq, group.group_id, 0, packet + sizeof(struct TOPOMSG_HEADER), sizeof(packet) - sizeof(struct TOPOMSG_HEADER), 0);

Loading…
Cancel
Save