Browse Source

Own recovery attempts in group sessions and require restored routes

proxy
evgeny 4 days ago
parent
commit
52ef7a4b3c
  1. 170
      src/routing_layer/topo_group.c
  2. 27
      src/routing_layer/topo_group.h
  3. 366
      src/routing_layer/topo_recovery.c
  4. 91
      src/routing_layer/topo_recovery.h
  5. 5
      tests/Makefile.am
  6. 11
      tests/test_etcp_lifecycle.c
  7. 14
      tests/test_group_ownership.c
  8. 196
      tests/test_group_recovery.c

170
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) { 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) {
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 && item->conn->peer_node_id == peer_id) return item; if (item->node_id == peer_id) return item;
} }
return NULL; return NULL;
} }
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->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) { 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", 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)group->group_id, (unsigned long long)peer->conn->peer_node_id,
(unsigned long long)peer->exchange_id, peer->table_sent, peer->table_received, (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; struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (peer->conn != conn) continue; if (peer->conn != conn) continue;
peer->exchange_id = 0; peer->table_received = 0; peer->table_sent = 0; 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", 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); (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); 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++) { for (int i = 0; i < chn; i++) {
struct TOPO_GROUP* cg = topo_groups_find(groups, chs[i]); struct TOPO_GROUP* cg = topo_groups_find(groups, chs[i]);
if (cg && cg->group_type == TOPO_GROUP_TYPE_CHAT) 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); u_free(chs);
@ -460,8 +562,9 @@ static struct TOPO_GROUP* topo_group_create(struct UTUN_INSTANCE* instance, uint
/* Полностью разрушает группу: recovery, connect, broadcast, conn_mgr, узлы. */ /* Полностью разрушает группу: recovery, connect, broadcast, conn_mgr, узлы. */
static void topo_group_destroy(struct TOPO_GROUP* group) { static void topo_group_destroy(struct TOPO_GROUP* group) {
if (!group) return; 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", 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); topo_recovery_cancel_all(group);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[GRP_DESTROY] 1 recovery_cancel done grp=%p", 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; 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", 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); (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_entry_free(e);
} }
queue_free(group->senders_list); 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, запрос таблицы. */ /* Добавляет пира в группу при 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 || !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) { 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 (!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)) { 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); (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id);
return -1; 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 стреляет ETCP_CONN_STATUS_UP дважды (UDP-линк, затем TCP-линк).
* Если conn уже в senders_list — повторно не обрабатываем, иначе active_conn_count * Если 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, (void*)conn, (unsigned long long)conn->peer_node_id, conn->state,
(unsigned long long)group->group_id, group->channel_id); (unsigned long long)group->group_id, group->channel_id);
if (topo_group_add_to_senders(group, conn) != 0) return -1; 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); 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; 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. */ /* Обрабатывает DOWN пира: чистит пути, каскадно удаляет узлы, шлёт withdraw. */
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason) { 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 (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; }
if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid instance"); 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; bool found_in_list = false;
struct ll_entry* e = group->senders_list->head; 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; struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)node_entry;
{ {
uint64_t next_hop = nq->node_id; uint64_t next_hop = nq->node_id;
struct ll_entry* pe = nq->paths ? nq->paths->head : NULL; uint16_t lost_rtt = UINT16_MAX;
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; } 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) { if (topo_group_remove_path(nq, conn) == 1) {
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, nq, next_hop); cascaded++; topo_recovery_add_node(group, key, next_hop, lost_rtt); cascaded++;
} }
nq->conn_presence = 0; nq->conn_up = 0; nq->conn_presence = 0; nq->conn_up = 0;
topo_fire_nodeinfo_cbk(conn->instance, group, nq); 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) { while (e) {
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) {
struct NODE_CONN_DIRECT* h = item->handle;
item->handle = NULL;
queue_remove_data(group->senders_list, e); queue_remove_data(group->senders_list, e);
topo_group_peer_free(item);
queue_entry_free(e); 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", 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, (void*)conn, (unsigned long long)conn->peer_node_id,
(unsigned long long)group->group_id, group->channel_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; break;
} }
e = e->next; e = e->next;
@ -1184,6 +1308,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
} }
/* Создание линков может вызвать транспортные callbacks: структура группы уже обновлена. */ /* Создание линков может вызвать транспортные callbacks: структура группы уже обновлена. */
topo_group_peer_progressed(group, from->peer_node_id);
node_conn_direct_update_node(group->instance, node_id); node_conn_direct_update_node(group->instance, node_id);
return 0; return 0;
@ -1263,12 +1388,21 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n
/* Добавляет conn в senders_list (дедуп), если его там ещё нет. */ /* Добавляет conn в senders_list (дедуп), если его там ещё нет. */
static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn || !group->senders_list) return -1; 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) 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; 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)); 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; } 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; 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) { 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); 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; 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); 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); 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); topo_group_send_table_request(group, conn);
} }

27
src/routing_layer/topo_group.h

@ -50,6 +50,7 @@ struct NAT_DETECTION;
struct CONN_MGR; struct CONN_MGR;
struct TOPO_RECOVERY_CTX; struct TOPO_RECOVERY_CTX;
struct TOPO_GROUP_CONNECT; struct TOPO_GROUP_CONNECT;
struct TOPO_PEER_REQUEST;
struct broadcast_ctx; struct broadcast_ctx;
struct radio_ctx; struct radio_ctx;
@ -180,8 +181,13 @@ struct TOPOMSG_ERR_GROUP_MISMATCH {
} __attribute__((packed)); } __attribute__((packed));
struct TOPO_GROUP_CONN_ITEM { 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 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 uint64_t exchange_id; // текущий исходящий REQUEST_TABLE
uint8_t table_received; // TABLE_COMPLETE для exchange_id принят uint8_t table_received; // TABLE_COMPLETE для exchange_id принят
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
@ -203,18 +209,19 @@ struct TOPO_GROUP {
uint64_t group_id; // уникальный идентификатор группы uint64_t group_id; // уникальный идентификатор группы
uint8_t group_type; // TOPO_GROUP_TYPE_* uint8_t group_type; // TOPO_GROUP_TYPE_*
struct UTUN_INSTANCE* instance; 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 ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B)
struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе
uint8_t ed25519_public_key[SC_PUBKEY_SIZE]; uint8_t ed25519_public_key[SC_PUBKEY_SIZE];
char channel_id[64]; // channel_id для групп типа CHAT char channel_id[64]; // channel_id для групп типа CHAT
struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы
struct TOPO_RECOVERY_CTX* recovery_list; // список активных процедур восстановления struct TOPO_RECOVERY_CTX* recovery; // один последовательный recovery на группу
struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT) struct TOPO_GROUP_CONNECT* connect; // авто-подключение к узлам группы (CHAT)
struct broadcast_ctx* broadcast; // broadcast protocol per-group context struct broadcast_ctx* broadcast; // broadcast protocol per-group context
struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context struct radio_ctx* radio; // radio (PTT walkie-talkie) per-group context
uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO) uint8_t radio_active; // 1 = я слушаю рацию этого канала (TOPO_FLAG_RADIO в NODEINFO)
struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов struct topo_node_cbk_entry* node_cbks; // цепочка подписчиков на события узлов
uint8_t stopping;
}; };
/* Отложенный REQUEST_TABLE: пир запросил таблицу группы, которой у нас ещё нет. /* Отложенный REQUEST_TABLE: пир запросил таблицу группы, которой у нас ещё нет.
@ -311,6 +318,20 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
* Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */ * Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */
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);
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. * @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable.
* *

366
src/routing_layer/topo_recovery.c

@ -1,16 +1,4 @@
/** #include <limits.h>
* @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 <stdlib.h>
#include <string.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"
@ -19,232 +7,188 @@
#include "topo_node.h" #include "topo_node.h"
#include "topo_group.h" #include "topo_group.h"
#include "topo_recovery.h" #include "topo_recovery.h"
#include "../transport_layer/node_conn_direct.h"
#include "../transport_layer/etcp.h" #include "../transport_layer/etcp.h"
#define RECOVERY_CAPACITY_INIT 16 struct recovery_node {
uint64_t node_id;
/* Выделение и инициализация контекста восстановления */ uint16_t rtt;
static struct TOPO_RECOVERY_CTX* topo_recovery_ctx_alloc(struct TOPO_GROUP* group, uint64_t next_hop_id) { uint8_t priority, tried, restored;
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; struct TOPO_RECOVERY_CTX {
ctx->capacity = RECOVERY_CAPACITY_INIT; struct TOPO_GROUP* group;
ctx->nodes = u_calloc(ctx->capacity, sizeof(struct TOPO_RECOVERY_NODE)); struct recovery_node* nodes;
if (!ctx->nodes) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery: nodes allocation failed"); u_free(ctx); return NULL; } size_t count, capacity;
return ctx; struct TOPO_PEER_REQUEST* request;
} uint64_t current_node, started_at;
enum topo_peer_phase phase;
static void topo_recovery_remove_candidate(struct TOPO_RECOVERY_CTX* ctx, size_t index) { void* wake;
uint64_t id = ctx->nodes[index].node_id; void* timer;
ctx->nodes[index] = ctx->nodes[--ctx->count]; unsigned attempts;
topo_node_registry_unref(ctx->instance->topo_groups, id); uint8_t started;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "recovery: release candidate=%016llx remaining=%zu", (unsigned long long)id, ctx->count); };
}
static void recovery_step(void* arg);
static void topo_recovery_ctx_free(struct TOPO_RECOVERY_CTX* ctx) {
if (!ctx) return; static void recovery_finish(struct TOPO_RECOVERY_CTX* ctx) {
while (ctx->count) topo_recovery_remove_candidate(ctx, ctx->count - 1); struct TOPO_GROUP* group = ctx->group;
u_free(ctx->nodes); group->recovery = NULL;
u_free(ctx); 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 int recovery_has_route(struct TOPO_GROUP* group, uint64_t node_id) {
static struct TOPO_RECOVERY_CTX* topo_recovery_find_by_next_hop(struct TOPO_GROUP* group, uint64_t next_hop_id) { struct TOPO_GROUP_NODE* node = topo_node_find_by_id(group, node_id);
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list; for (struct ll_entry* e = node && node->paths ? node->paths->head : NULL; e; e = e->next) {
while (ctx) { if (!ctx->started && ctx->next_hop_id == next_hop_id) return ctx; ctx = ctx->next; } struct ETCP_CONN* conn = ((struct TOPO_NODEPATH*)e)->conn;
return NULL; if (conn && conn->links_up && !conn->close_requested && topo_group_peer_ready(group, conn->peer_node_id)) return 1;
} }
return 0;
/* Вынимает контекст из связного списка group->recovery_list */
static void topo_recovery_ctx_remove(struct TOPO_GROUP* group, struct TOPO_RECOVERY_CTX* ctx) {
struct TOPO_RECOVERY_CTX** pp = &group->recovery_list;
while (*pp) { if (*pp == ctx) { *pp = ctx->next; return; } pp = &(*pp)->next; }
} }
static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx); static size_t recovery_remaining(struct TOPO_RECOVERY_CTX* ctx) {
size_t remaining = 0;
/* Коллбэк ncd: подключение удалось или провалилось */ for (size_t i = 0; i < ctx->count; i++) {
static void topo_recovery_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { struct recovery_node* node = &ctx->nodes[i];
struct TOPO_RECOVERY_CTX* ctx = (struct TOPO_RECOVERY_CTX*)arg; int restored = recovery_has_route(ctx->group, node->node_id);
struct ETCP_CONN* conn = node_conn_direct_get_conn(h); if (restored != node->restored) {
uint64_t node_id = node_conn_direct_node_id(h); node->restored = (uint8_t)restored;
if (node_id != ctx->current_node_id) return; DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery route: group=%016llx node=%016llx restored=%d",
if (ctx->connect_timer) { uasync_cancel_timeout(ctx->instance->ua, ctx->connect_timer); ctx->connect_timer = NULL; } (unsigned long long)ctx->group->group_id, (unsigned long long)node->node_id, restored);
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;
} }
DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: BGP ownership transfer failed node=%016llx", (unsigned long long)node_id); if (!restored) remaining++;
ctx->next = ctx->group->recovery_list; ctx->group->recovery_list = ctx;
} }
node_conn_direct_close(h); return remaining;
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);
} }
/* Таймаут 2с: закрывает handle, удаляет узел, переходит к следующему */ static struct recovery_node* recovery_best(struct TOPO_RECOVERY_CTX* ctx) {
static void topo_recovery_timeout_cb(void* arg) { struct recovery_node* best = NULL;
struct TOPO_RECOVERY_CTX* ctx = (struct TOPO_RECOVERY_CTX*)arg; for (size_t i = 0; i < ctx->count; i++) {
ctx->connect_timer = NULL; struct recovery_node* node = &ctx->nodes[i];
uint64_t node_id = ctx->current_node_id; if (node->tried || node->restored) continue;
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 (!best || node->priority > best->priority || (node->priority == best->priority &&
if (ctx->current_handle) { (node->rtt < best->rtt || (node->rtt == best->rtt && node->node_id < best->node_id)))) best = node;
node_conn_direct_close(ctx->current_handle);
ctx->current_handle = NULL;
} }
for (size_t i = 0; i < ctx->count; i++) return best;
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);
} }
/* Выбирает узел: сперва is_next_hop с мин. RTT, за ним — остальные по мин. RTT */ static void recovery_timeout(void* arg) {
static struct TOPO_RECOVERY_NODE* topo_recovery_find_best(struct TOPO_RECOVERY_CTX* ctx) { struct TOPO_RECOVERY_CTX* ctx = arg;
struct TOPO_RECOVERY_NODE* best = NULL; ctx->timer = NULL;
for (size_t i = 0; i < ctx->count; i++) recovery_step(ctx);
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 topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx) { static void recovery_wake(void* arg) {
while (1) { struct TOPO_RECOVERY_CTX* ctx = arg;
if (ctx->count == 0) { ctx->wake = NULL;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx all nodes exhausted, no reconnections", (unsigned long long)ctx->next_hop_id); recovery_step(ctx);
topo_recovery_ctx_remove(ctx->group, ctx);
topo_recovery_ctx_free(ctx);
return;
}
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;
}
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;
}
ctx->connect_timer = uasync_set_timeout(ctx->instance->ua, TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10,
ctx, topo_recovery_timeout_cb, "recovery_timeout");
return;
}
} }
void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t next_hop) { static void recovery_step(void* arg) {
if (!group || !nq) return; struct TOPO_RECOVERY_CTX* ctx = arg;
struct TOPO_NODE* ni = topo_node_registry_find(group->instance->topo_groups, nq->node_id); struct TOPO_GROUP* group = ctx->group;
if (!ni) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: missing candidate=%016llx", (unsigned long long)nq->node_id); return; } if (ctx->timer) { uasync_cancel_timeout(group->instance->ua, ctx->timer); ctx->timer = NULL; }
uint64_t node_id = nq->node_id; 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);
}
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 (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;
}
}
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);
}
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) { if (!ctx) {
ctx = topo_recovery_ctx_alloc(group, next_hop); ctx = u_calloc(1, sizeof(*ctx));
if (!ctx) return; if (!ctx) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery context allocation failed"); return; }
ctx->next = group->recovery_list; ctx->group = group; group->recovery = ctx;
group->recovery_list = ctx; }
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) {
for (size_t i = 0; i < ctx->count; i++) if (ctx->nodes[i].node_id == node_id) return; size_t capacity = ctx->capacity ? ctx->capacity * 2 : 16;
uint16_t rtt = topo_get_chain_rtt(nq); struct recovery_node* nodes = u_realloc(ctx->nodes, capacity * sizeof(*nodes));
if (!nodes) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery candidate allocation failed"); return; }
if (ctx->count >= ctx->capacity) { ctx->nodes = nodes; ctx->capacity = 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;
} }
topo_node_registry_ref(group->instance->topo_groups, node_id); topo_node_registry_ref(group->instance->topo_groups, node_id);
ctx->nodes[ctx->count].node_id = node_id; ctx->nodes[ctx->count++] = (struct recovery_node){ .node_id = node_id, .rtt = rtt, .priority = node_id == next_hop };
ctx->nodes[ctx->count].min_rtt = rtt; DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery target: group=%016llx node=%016llx next=%016llx rtt=%u targets=%zu",
ctx->nodes[ctx->count].is_next_hop = (uint8_t)(node_id == next_hop ? 1 : 0); (unsigned long long)group->group_id, (unsigned long long)node_id, (unsigned long long)next_hop, rtt, ctx->count);
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);
} }
void topo_recovery_start(struct TOPO_GROUP* group) { void topo_recovery_start(struct TOPO_GROUP* group) {
if (!group) return; if (!group || !group->recovery || group->stopping) return;
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list; group->recovery->started = 1;
int started_ctxs = 0; topo_recovery_changed(group);
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");
}
} }
void topo_recovery_cancel_for_node(struct TOPO_GROUP* group, uint64_t node_id) { void topo_recovery_changed(struct TOPO_GROUP* group) {
if (!group) return; struct TOPO_RECOVERY_CTX* ctx = group ? group->recovery : NULL;
struct TOPO_RECOVERY_CTX** pp = &group->recovery_list; if (!ctx || !ctx->started || group->stopping || ctx->wake) return;
while (*pp) { ctx->wake = uasync_call_soon(group->instance->ua, ctx, recovery_wake);
struct TOPO_RECOVERY_CTX* ctx = *pp; if (!ctx->wake) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "recovery wake allocation failed group=%016llx", (unsigned long long)group->group_id);
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_cancel_all(struct TOPO_GROUP* group) { void topo_recovery_cancel_all(struct TOPO_GROUP* group) {
if (!group) return; if (!group || !group->recovery) return;
struct TOPO_RECOVERY_CTX* ctx = group->recovery_list; DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery cancelled group=%016llx", (unsigned long long)group->group_id);
while (ctx) { recovery_finish(group->recovery);
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;
} }

91
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
*
* Группировка: <мы> -> <failed_peer N> -> <next_hop A, B...>
* Каждый next_hop со своим subtree — отдельный recovery-контекст.
* Восстановление next_hop автоматически оживляет его subtree через BGP.
*
* Отмена: topo_group_new_conn (узел появился в сети) сканирует все active recovery,
* при совпадении отменяет контекст целиком. Самоуничтожение при исчерпании списков.
*/
#ifndef TOPO_RECOVERY_H #ifndef TOPO_RECOVERY_H
#define TOPO_RECOVERY_H #define TOPO_RECOVERY_H
#include <stdint.h>
#ifdef __cplusplus #ifdef __cplusplus
extern "C" { extern "C" {
#endif #endif
#include <stdint.h>
#include <stddef.h>
struct TOPO_GROUP; struct TOPO_GROUP;
struct TOPO_GROUP_NODE; struct TOPO_RECOVERY_CTX;
struct NODE_CONN_DIRECT;
#define TOPO_RECOVERY_CONNECT_TIMEOUT_MS 2000 #define TOPO_RECOVERY_CONNECT_TIMEOUT_MS 2000
#define TOPO_RECOVERY_SYNC_TIMEOUT_MS 5000
struct TOPO_RECOVERY_NODE {
uint64_t node_id; /* Один recovery-контекст на группу, одна текущая попытка присоединения.
uint16_t min_rtt; /* 0.1ms, из connectivity-проб (interface/nat/real) */ * add_node сохраняет цель и кандидатуру до удаления registry ref;
uint8_t is_next_hop; /* 1 = прямой downstream отвалившегося узла, пробуется первым */ * start вызывается после удаления всех путей через потерянного пира.
}; * Приоритет: бывшие downstream next-hop, затем минимальный RTT, затем node_id.
* Кандидат проверяется не более одного раза за цикл. Транспорт принадлежит
struct TOPO_RECOVERY_CTX { * групповой сессии: recovery владеет только отменяемым TOPO_PEER_REQUEST.
struct TOPO_RECOVERY_CTX* next; /* следующий в group->recovery_list */ *
struct UTUN_INSTANCE* instance; * Успех — все потерянные узлы снова имеют живой путь через READY-пира этой группы.
struct TOPO_GROUP* group; * UP и даже READY без нужных маршрутов не завершают recovery. Частичный результат
uint64_t next_hop_id; /* ключ группировки: node_id next-hop'а от failed_peer */ * сохраняет полезное присоединение и продолжает оставшиеся цели. Маршрут через
struct TOPO_RECOVERY_NODE* nodes; /* кандидаты на восстановление (прямые) */ * другого READY-пира также закрывает цель. При исчерпании кандидатов оставшиеся
size_t count; /* текущее количество */ * цели логируются, контекст завершается без скрытого повторного цикла.
size_t capacity; /* выделенная ёмкость */ *
uint64_t current_node_id; /* node_id в текущей попытке, 0=нет активной */ * CONNECTING ограничен отдельным таймаутом; SYNCING — временем без прогресса.
struct NODE_CONN_DIRECT* current_handle; /* handle текущей попытки node_conn_direct_open */ * NODEINFO/смена состояния лишь планируют проверку через call_soon: текущий
void* connect_timer; /* внешний таймер 2с (uasync) */ * BGP callback завершается до изменения попытки или освобождения контекста.
uint8_t started; /* 0=сбор узлов (add_node), 1=перебор запущен */ * Остановка группы отменяет recovery до освобождения сессий и таблицы узлов. */
}; void topo_recovery_add_node(struct TOPO_GROUP* group, uint64_t node_id, uint64_t next_hop, uint16_t rtt);
/**
* Добавляет узел в 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 если есть каскадные узлы.
*/
void topo_recovery_start(struct TOPO_GROUP* group); void topo_recovery_start(struct TOPO_GROUP* group);
void topo_recovery_changed(struct TOPO_GROUP* group);
/**
* Сканирует все active recovery: если node_id найден в любом контексте —
* отменяет текущую попытку connect и освобождает контекст.
* Вызывается из topo_group_new_conn (узел появился в сети).
*/
void topo_recovery_cancel_for_node(struct TOPO_GROUP* group, uint64_t node_id);
/** Отменяет все recovery-контексты. Вызывается из topo_group_destroy. */
void topo_recovery_cancel_all(struct TOPO_GROUP* group); void topo_recovery_cancel_all(struct TOPO_GROUP* group);
#ifdef __cplusplus #ifdef __cplusplus
} }
#endif #endif
#endif #endif

5
tests/Makefile.am

@ -66,6 +66,7 @@ check_PROGRAMS = \
test_node_conn_direct \ test_node_conn_direct \
test_group_ownership \ test_group_ownership \
test_group_exchange \ test_group_exchange \
test_group_recovery \
test_node_snapshot \ test_node_snapshot \
test_ncd_config \ test_ncd_config \
test_db_sync \ 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_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_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_SOURCES = test_node_snapshot.c
test_node_snapshot_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_node_snapshot_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_node_snapshot_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_node_snapshot_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

11
tests/test_etcp_lifecycle.c

@ -111,15 +111,14 @@ static void recovery_cases(struct UTUN_INSTANCE* inst) {
for (int scenario=0;scenario<3;scenario++) { for (int scenario=0;scenario<3;scenario++) {
struct TOPO_NODE* ni = u_calloc(1,sizeof(*ni)); CHECK(ni); ni->node_id = node.node_id; 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); 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); 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_list->count == 1); CHECK(ni->group_ref_count == 2 && group.recovery);
if (scenario == 1) { topo_recovery_add_node(&group,&node,45); CHECK(ni->group_ref_count == 3); } 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); 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); if (scenario < 2) topo_recovery_cancel_all(&group);
else if (scenario == 1) topo_recovery_cancel_all(&group);
else { topo_recovery_start(&group); topo_recovery_cancel_all(&group); } else { topo_recovery_start(&group); topo_recovery_cancel_all(&group); }
uasync_poll(inst->ua,0); 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); printf("PASS recovery ownership scenario %d\n",scenario+1);
} }
queue_free(groups.node_registry); inst->topo_groups = NULL; queue_free(groups.node_registry); inst->topo_groups = NULL;

14
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); assert(queue_data_put_with_index(group->nodes, entry) == 0);
topo_group_remove_conn(group, conn, reason); topo_group_remove_conn(group, conn, reason);
assert(!topo_node_find_by_id(group, id)); 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); assert(!conn->close_requested);
} }
@ -51,10 +51,22 @@ int main(void) {
/* Создание группы при уже существующем транспорте читает payload записи conn. */ /* Создание группы при уже существующем транспорте читает payload записи conn. */
struct TOPO_GROUP* c = topo_groups_create_group(inst->topo_groups, 44, TOPO_GROUP_TYPE_UTUN, NULL); assert(c); 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); 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_LOCAL_LEAVE);
check_deliberate_leave(c, conn, TOPO_REMOVE_REMOTE_LEAVE); check_deliberate_leave(c, conn, TOPO_REMOVE_REMOTE_LEAVE);
check_deliberate_leave(c, conn, TOPO_REMOVE_MEMBER_INVALID); 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); 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);
assert(topo_group_new_conn(a, conn) == 0 && queue_entry_count(a->senders_list) == 1); assert(topo_group_new_conn(a, conn) == 0 && queue_entry_count(a->senders_list) == 1);

196
tests/test_group_recovery.c

@ -0,0 +1,196 @@
#include <assert.h>
#include <string.h>
#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;
}
Loading…
Cancel
Save