Browse Source

Prune unconfirmed BGP paths and restart stalled exchanges

proxy
evgeny 3 days ago
parent
commit
39ce249361
  1. 77
      src/routing_layer/topo_group.c
  2. 9
      src/routing_layer/topo_group.h
  3. 5
      tests/test_bgp_paths.c
  4. 23
      tests/test_group_exchange.c
  5. 2
      tests/test_group_recovery.c

77
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_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); 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 topo_advert {
struct ll_entry ll; 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) { 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); 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) { static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) {
topo_exchange_cancel(peer);
topo_tx_clear(peer); topo_tx_clear(peer);
topo_adverts_clear(peer); topo_adverts_clear(peer);
while (peer->requests) { 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) { 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)) { if (topo_group_peer_ready(group, peer->node_id)) {
topo_exchange_cancel(peer);
peer->retained = 1; topo_group_connect_on_up(group, peer->conn); peer->retained = 1; topo_group_connect_on_up(group, peer->conn);
} }
topo_group_peer_progressed(group, peer->node_id); 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); topo_group_peer_progressed(peer->group, peer->node_id);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group session started: group=%016llx peer=%016llx local=%016llx", 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); (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) { 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)); 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) { static void topo_group_send_resync(struct ETCP_CONN* conn) {
struct TOPOMSG_RESYNC msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_RESYNC }; struct TOPOMSG_RESYNC msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_RESYNC };
topo_send_control(conn, &msg, sizeof(msg)); 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", 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, 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); (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_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_REQUEST_TABLE) topo_group_handle_request_table(group, from_conn);
else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE && !peer->table_received) { 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); peer->table_received = 1; topo_group_log_exchange(group, peer);
} else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) { } else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) {
const struct TOPOMSG_ERR_GROUP_MISMATCH* err = (const struct TOPOMSG_ERR_GROUP_MISMATCH*)data; 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; 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_exchange_cancel(item);
topo_tx_clear(item); topo_adverts_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);
@ -1113,6 +1155,29 @@ static void topo_group_routes_changed(struct TOPO_GROUP* group, struct TOPO_GROU
topo_recovery_changed(group); 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 ===== // ===== NODEINFO process =====
/* Шлёт пиру ошибку несоответствия типа группы (ERR_GROUP_MISMATCH). */ /* Шлёт пиру ошибку несоответствия типа группы (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); 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_GROUP_CONN_ITEM* peer = topo_group_peer(group, from->peer_node_id);
((struct TOPO_NODEPATH*)nodeinfo1->paths->tail)->exchange = peer ? peer->local_epoch : 0; ((struct TOPO_NODEPATH*)nodeinfo1->paths->tail)->exchange = peer ? peer->local_epoch : 0;
if (node_id != group->instance->node_id) { 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", 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); (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_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0;
topo_exchange_arm(peer);
topo_recovery_changed(peer->group); topo_recovery_changed(peer->group);
topo_group_connect_changed(peer->group); topo_group_connect_changed(peer->group);
} }

9
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/leave/stop отменяют остаток и waiter; уже
* переданные транспорту пакеты защищены поколениями на принимающей стороне. * переданные транспорту пакеты защищены поколениями на принимающей стороне.
* *
* Новое присоединение и REINIT создают новое поколение; REINIT сохраняет пути. * Новое присоединение и REINIT создают новое поколение; REINIT временно сохраняет
* пути. Каждый NODEINFO помечает путь текущим local_epoch. На TABLE_COMPLETE
* удаляются неподтверждённые пути этого соседа; альтернативы остаются.
* Если обмен не продвигается TOPO_GROUP_SYNC_TIMEOUT_MS, все пути этого соседа
* удаляются и JOIN начинается заново. Ошибка применения/отправки запрещает READY.
* Поколения защищают сообщения, а доступность существующего пути определяется * Поколения защищают сообщения, а доступность существующего пути определяется
* его соединением. Recovery дополнительно ожидает READY группового пира. * его соединением. Recovery дополнительно ожидает READY группового пира.
* JOIN_READY чата означает добавление мембера и не заменяет этот обмен. * JOIN_READY чата означает добавление мембера и не заменяет этот обмен.
@ -181,6 +185,7 @@ struct TOPO_GROUP_CONN_ITEM {
struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle)
struct TOPO_PEER_REQUEST* requests; // запросы инициаторов; каждый закрывает свой struct TOPO_PEER_REQUEST* requests; // запросы инициаторов; каждый закрывает свой
uint64_t progress; // время последнего прогресса обмена (timebase) uint64_t progress; // время последнего прогресса обмена (timebase)
void* exchange_timer; // watchdog незавершённого обмена; отменяется при READY/DOWN/free
uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint8_t retained; // участие запрошено постоянным владельцем или достигло READY
uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения
uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch
@ -193,6 +198,8 @@ struct TOPO_GROUP_CONN_ITEM {
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
}; };
#define TOPO_GROUP_SYNC_TIMEOUT_MS 5000 /* отсутствие прогресса: удалить пути пира и начать новый JOIN */
struct ETCP_LINK; struct ETCP_LINK;
#define TOPO_NODE_REGISTRY_HASH_SIZE 256 #define TOPO_NODE_REGISTRY_HASH_SIZE 256

5
tests/test_bgp_paths.c

@ -144,6 +144,11 @@ int main(void) {
create(); create();
unsigned triangle = edge(0, 1) | edge(1, 2) | edge(0, 2) | edge(2, 3); unsigned triangle = edge(0, 1) | edge(1, 2) | edge(0, 2) | edge(2, 3);
change(triangle); settle(); verify(); 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 & ~edge(0, 2)); settle(); verify();
change(triangle); settle(); verify(); change(triangle); settle(); verify();
change(triangle & ~edge(1, 2) & ~edge(0, 2)); settle(); verify(); change(triangle & ~edge(1, 2) & ~edge(0, 2)); settle(); verify();

23
tests/test_group_exchange.c

@ -8,6 +8,7 @@
#include "transport_layer/node_conn_direct.h" #include "transport_layer/node_conn_direct.h"
#include "../lib/u_async.h" #include "../lib/u_async.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/platform_compat.h"
static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) { static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) {
assert(group->senders_list->head && !group->senders_list->head->next); 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)->table_sent);
assert(peer(a)->tx_failed); assert(peer(a)->tx_failed);
conn->send_input_q->size_limit = saved_limit; 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); accept_join(a, conn);
nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 1);
complete(a, conn, peer(a)->local_epoch); request(a, conn); complete(a, conn, peer(a)->local_epoch); request(a, conn);
assert(topo_group_peer_ready(a, node.node_id)); 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)); 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));
assert(!topo_node_find_by_id(a, node.node_id)); 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; a_id = peer(a)->local_epoch;
topo_group_remove_conn(a, conn, TOPO_REMOVE_LOCAL_LEAVE); topo_group_remove_conn(a, conn, TOPO_REMOVE_LOCAL_LEAVE);
complete(a, conn, a_id); complete(a, conn, a_id);

2
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); topo_node_registry_ref(f->inst->topo_groups, node->node_id);
uint64_t hops[] = { node->node_id, f->ids[via] }; 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(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); assert(queue_data_put_with_index(f->group->nodes, entry) == 0);
topo_recovery_changed(f->group); topo_recovery_changed(f->group);
} }

Loading…
Cancel
Save