Browse Source

Serialize group topology lazily through a cancellable backpressure sender

proxy
evgeny 4 days ago
parent
commit
0756a6e020
  1. 269
      src/routing_layer/topo_group.c
  2. 15
      src/routing_layer/topo_group.h
  3. 2
      src/routing_layer/topo_node.c
  4. 2
      src/transport_layer/etcp_connections.c
  5. 4
      tests/test_etcp_connect.c
  6. 16
      tests/test_group_exchange.c
  7. 73
      tests/test_group_recovery.c

269
src/routing_layer/topo_group.c

@ -44,6 +44,11 @@ struct TOPO_PEER_REQUEST {
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 void topo_tx_clear(struct TOPO_GROUP_CONN_ITEM* peer);
static int topo_tx_control(struct TOPO_GROUP_CONN_ITEM* peer, const void* data, size_t size);
static int topo_tx_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t node_id);
static int topo_tx_table(struct TOPO_GROUP_CONN_ITEM* peer);
static int topo_control_wait(struct ETCP_CONN* conn, const void* data, size_t size);
static 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) {
@ -59,6 +64,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) {
topo_tx_clear(peer);
while (peer->requests) {
struct TOPO_PEER_REQUEST* request = peer->requests;
peer->requests = request->next; request->peer = NULL; request->next = NULL;
@ -119,7 +125,7 @@ int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO
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 || peer->tx_failed) 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;
}
@ -163,7 +169,7 @@ static struct TOPOMSG_HEADER topo_header(struct TOPO_GROUP_CONN_ITEM* peer, uint
.src_epoch = peer->local_epoch, .dst_epoch = peer->peer_epoch };
}
static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t size) {
static int topo_emit_control(struct ETCP_CONN* conn, const void* data, size_t size) {
struct ll_entry* e = ll_alloc_lldgram(size);
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group control allocation failed size=%zu", size); return -1; }
memcpy(e->dgram, data, size); e->len = size;
@ -173,7 +179,20 @@ static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t si
queue_dgram_free(e); queue_entry_free(e); return -1;
}
/* Сообщения сессии проходят одну FIFO. Stateless RESYNC/REJECT также ждут
* пустой транспортной очереди, но не создают групповое участие. */
static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t size) {
if (size >= sizeof(struct TOPOMSG_HEADER)) {
const struct TOPOMSG_HEADER* h = data;
struct TOPO_GROUP* group = topo_groups_find(conn->instance->topo_groups, h->group_id);
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id);
if (peer && peer->conn == conn && peer->local_epoch == h->src_epoch) return topo_tx_control(peer, data, size);
}
return topo_control_wait(conn, data, size);
}
static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) {
topo_tx_clear(peer); peer->tx_failed = 0;
peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0;
uint64_t epoch;
if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) {
@ -212,7 +231,7 @@ static int topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id);
if (!peer || !peer->accepted) return -1;
struct TOPOMSG_HEADER msg = topo_header(peer, TOPO_SUBCMD_TABLE_COMPLETE);
return topo_send_control(conn, &msg, sizeof(msg));
return topo_emit_control(conn, &msg, sizeof(msg));
}
static void topo_group_reject(struct ETCP_CONN* conn, const struct TOPOMSG_HEADER* h, enum topo_join_reject reason) {
@ -225,7 +244,6 @@ static void topo_group_reject(struct ETCP_CONN* conn, const struct TOPOMSG_HEADE
static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t target_id);
static int topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn, const struct TOPOMSG_JOIN_GROUP* msg);
static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn);
@ -660,7 +678,7 @@ void topo_group_on_activity_change(struct UTUN_INSTANCE* instance, int active, v
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn);
se = se->next;
}
}
@ -806,7 +824,7 @@ void topo_group_set_radio(struct TOPO_GROUP* group, int on) {
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn && group->local_node)
topo_group_send_nodeinfo(group, group->local_node, item->conn, 0);
topo_group_send_nodeinfo(group, group->local_node, item->conn);
se = se->next;
}
}
@ -1235,8 +1253,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
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) {
uint16_t out_rtt = topo_get_chain_rtt(nodeinfo1);
topo_group_send_nodeinfo(group, nodeinfo1, item->conn, out_rtt);
topo_group_send_nodeinfo(group, nodeinfo1, item->conn);
}
}
}
@ -1289,7 +1306,7 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
}
/* Сериализует и шлёт NODEINFO узла конкретному пиру. */
int topo_group_send_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; }
size_t max_sz = sizeof(struct TOPOMSG_NODEINFO_PKT) + TOPO_NODE_WIRE_MAX_SIZE;
@ -1334,6 +1351,206 @@ int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* n
return 0;
}
/* NODEINFO сериализуется только при отправке: указатели на узлы не переживают callback. */
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn) {
if (!group || !node || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: invalid arguments"); return -1; }
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id);
if (!peer || !peer->accepted) return 0;
return topo_tx_node(peer, node->node_id);
}
enum topo_tx_kind { TOPO_TX_CONTROL, TOPO_TX_NODE, TOPO_TX_TABLE };
struct topo_tx_item {
struct topo_tx_item* next;
enum topo_tx_kind kind;
uint64_t node_id;
struct ll_entry* packet;
uint64_t* nodes;
size_t count, pos, sent;
};
static void topo_tx_item_free(struct topo_tx_item* item) {
if (item->packet) { queue_dgram_free(item->packet); queue_entry_free(item->packet); }
u_free(item->nodes); u_free(item);
}
static void topo_tx_clear(struct TOPO_GROUP_CONN_ITEM* peer) {
if (peer->tx_wake) { uasync_call_soon_cancel(peer->group->instance->ua, peer->tx_wake); peer->tx_wake = NULL; }
if (peer->conn && peer->conn->send_input_q) queue_waiter_cancel(peer->conn->send_input_q, &peer->tx_waiter);
while (peer->tx_head) {
struct topo_tx_item* item = peer->tx_head; peer->tx_head = item->next; topo_tx_item_free(item);
}
peer->tx_tail = NULL;
}
static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group sender failed: group=%016llx peer=%016llx epoch=%016llx",
(unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch);
topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0;
topo_recovery_changed(peer->group);
}
static void topo_tx_wait(void* arg);
static void topo_tx_schedule(struct TOPO_GROUP_CONN_ITEM* peer) {
if (!peer->tx_head || peer->tx_failed || peer->tx_wake || peer->tx_waiter.internal || peer->tx_waiter.call_soon_id) return;
peer->tx_wake = uasync_call_soon(peer->group->instance->ua, peer, topo_tx_wait);
if (!peer->tx_wake) topo_tx_fail(peer);
}
static int topo_tx_append(struct TOPO_GROUP_CONN_ITEM* peer, struct topo_tx_item* item) {
if (!item || peer->tx_failed) { if (item) topo_tx_item_free(item); topo_tx_fail(peer); return -1; }
if (peer->tx_tail) peer->tx_tail->next = item; else peer->tx_head = item;
peer->tx_tail = item; topo_tx_schedule(peer);
return peer->tx_failed ? -1 : 0;
}
static int topo_tx_control(struct TOPO_GROUP_CONN_ITEM* peer, const void* data, size_t size) {
struct topo_tx_item* item = u_calloc(1, sizeof(*item));
if (!item) { topo_tx_fail(peer); return -1; }
item->packet = ll_alloc_lldgram(size);
if (!item->packet) { u_free(item); topo_tx_fail(peer); return -1; }
memcpy(item->packet->dgram, data, size); item->packet->len = size;
item->kind = TOPO_TX_CONTROL;
return topo_tx_append(peer, item);
}
static int topo_tx_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t node_id) {
/* Последнее изменение занимает своё место после уже ожидающих WITHDRAW.
* На один узел хранится не более одного несериализованного NODEINFO. */
struct topo_tx_item** link = &peer->tx_head;
struct topo_tx_item* prev = NULL;
struct topo_tx_item* item = NULL;
while (*link) {
if ((*link)->kind == TOPO_TX_NODE && (*link)->node_id == node_id) {
item = *link; *link = item->next;
if (peer->tx_tail == item) peer->tx_tail = prev;
item->next = NULL; break;
}
prev = *link; link = &(*link)->next;
}
if (!item) item = u_calloc(1, sizeof(*item));
if (!item) { topo_tx_fail(peer); return -1; }
item->kind = TOPO_TX_NODE; item->node_id = node_id;
return topo_tx_append(peer, item);
}
static int topo_tx_table(struct TOPO_GROUP_CONN_ITEM* peer) {
for (struct topo_tx_item* item = peer->tx_head; item; item = item->next)
if (item->kind == TOPO_TX_TABLE) return 0;
struct TOPO_GROUP* group = peer->group;
struct topo_tx_item* item = u_calloc(1, sizeof(*item));
if (!item) { topo_tx_fail(peer); return -1; }
item->kind = TOPO_TX_TABLE;
size_t capacity = (size_t)queue_entry_count(group->nodes) + 1;
item->nodes = u_malloc(capacity * sizeof(*item->nodes));
if (!item->nodes) { u_free(item); topo_tx_fail(peer); return -1; }
item->nodes[item->count++] = group->local_node->node_id;
for (struct ll_entry* e = group->nodes->head; e; e = e->next) {
struct TOPO_GROUP_NODE* node = (struct TOPO_GROUP_NODE*)e;
if (node != group->local_node && topo_group_should_send_to(node, peer->node_id)) item->nodes[item->count++] = node->node_id;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group table queued: group=%016llx peer=%016llx epoch=%016llx nodes=%zu",
(unsigned long long)group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch, item->count);
return topo_tx_append(peer, item);
}
static int topo_tx_emit_node(struct TOPO_GROUP_CONN_ITEM* peer, uint64_t id) {
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);
if (!node || (node != group->local_node && !topo_group_should_send_to(node, peer->node_id))) return 0;
return topo_emit_nodeinfo(group, node, peer->conn, node == group->local_node ? 0 : topo_get_chain_rtt(node)) < 0 ? -1 : 1;
}
static void topo_tx_ready(struct ll_queue* queue, void* arg) {
struct TOPO_GROUP_CONN_ITEM* peer = arg;
/* Deferred waiters must recheck: another producer can occupy the queue first. */
if (queue_entry_count(queue)) { topo_tx_schedule(peer); return; }
struct topo_tx_item* item = peer->tx_head;
if (!item) return;
int result = 0, complete = 1;
if (item->kind == TOPO_TX_CONTROL) {
result = etcp_send(peer->conn, item->packet);
if (result == 0) item->packet = NULL;
} else if (item->kind == TOPO_TX_NODE) result = topo_tx_emit_node(peer, item->node_id);
else {
if (item->pos < item->count) {
result = topo_tx_emit_node(peer, item->nodes[item->pos++]);
if (result > 0) item->sent++;
complete = 0;
} else {
result = topo_group_send_table_complete(peer->group, peer->conn);
if (result == 0) {
peer->table_sent = 1; topo_group_log_exchange(peer->group, peer);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "group table sent: group=%016llx peer=%016llx epoch=%016llx nodes=%zu",
(unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id,
(unsigned long long)peer->local_epoch, item->sent);
}
}
}
if (result < 0) { topo_tx_fail(peer); return; }
topo_group_peer_progressed(peer->group, peer->node_id);
if (complete) {
peer->tx_head = item->next;
if (!peer->tx_head) peer->tx_tail = NULL;
topo_tx_item_free(item);
}
topo_tx_schedule(peer);
}
static void topo_tx_wait(void* arg) {
struct TOPO_GROUP_CONN_ITEM* peer = arg;
peer->tx_wake = NULL;
if (!peer->conn || !peer->conn->send_input_q || peer->conn->close_requested ||
queue_waiter_wait(peer->conn->send_input_q, &peer->tx_waiter, topo_tx_ready, peer) < 0) topo_tx_fail(peer);
}
/* Небольшой ответ вне групповой сессии владеет NCD только до постановки пакета
* в транспорт. CLOSED/DOWN отменяет waiter до уничтожения транспортной очереди. */
struct topo_control {
struct NODE_CONN_DIRECT* handle;
struct ll_queue* queue;
struct queue_waiter_handle waiter;
struct ll_entry* packet;
};
static void topo_control_free(struct topo_control* control) {
queue_waiter_cancel(control->queue, &control->waiter);
node_conn_direct_close(control->handle);
queue_dgram_free(control->packet); queue_entry_free(control->packet); u_free(control);
}
static void topo_control_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) {
(void)handle;
if (event == NCD_EVENT_UP) return;
DEBUG_WARN(DEBUG_CATEGORY_BGP, "pending group control cancelled transport_event=%d", event);
topo_control_free(arg);
}
static void topo_control_ready(struct ll_queue* queue, void* arg) {
struct topo_control* control = arg;
if (queue_entry_count(queue)) {
if (queue_waiter_wait(queue, &control->waiter, topo_control_ready, control) >= 0) return;
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot rearm group control waiter");
} else if (etcp_send(node_conn_direct_get_conn(control->handle), control->packet) == 0) control->packet = NULL;
else DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send pending group control");
topo_control_free(control);
}
static int topo_control_wait(struct ETCP_CONN* conn, const void* data, size_t size) {
if (!conn->send_input_q || conn->close_requested) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "group control on closing transport"); return -1; }
if (!queue_entry_count(conn->send_input_q)) return topo_emit_control(conn, data, size);
struct topo_control* control = u_calloc(1, sizeof(*control));
if (!control) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group control allocation failed"); return -1; }
control->queue = conn->send_input_q; control->packet = ll_alloc_lldgram(size);
if (!control->packet) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group control packet allocation failed"); u_free(control); return -1; }
memcpy(control->packet->dgram, data, size); control->packet->len = size;
if (node_conn_direct_open(conn->instance, conn->peer_node_id, topo_control_event, control, &control->handle, NULL) == NCD_ERR ||
queue_waiter_wait(control->queue, &control->waiter, topo_control_ready, control) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot retain pending group control"); topo_control_free(control); return -1;
}
return 0;
}
/* Добавляет conn в senders_list (дедуп), если его там ещё нет. */
static int topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn || !group->senders_list) return -1;
@ -1420,7 +1637,7 @@ static void topo_leave_ready(struct ll_queue* queue, void* arg) {
(unsigned long long)leave->group_id);
topo_leave_finish(leave); return;
}
int result = topo_send_control(conn, &leave->message, sizeof(leave->message));
int result = topo_emit_control(conn, &leave->message, sizeof(leave->message));
if (result != 0) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send LEAVE_GROUP group=%016llx", (unsigned long long)leave->group_id);
else DEBUG_INFO(DEBUG_CATEGORY_BGP, "LEAVE_GROUP sent: peer=%016llx group=%016llx",
(unsigned long long)node_conn_direct_node_id(leave->handle), (unsigned long long)leave->group_id);
@ -1454,38 +1671,10 @@ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* han
return 0;
}
/* Шлёт пиру полную таблицу узлов группы (все узлы, кроме достижимых через него). */
static int topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
uint64_t target = conn->peer_node_id;
struct ll_entry* e = group->nodes ? group->nodes->head : NULL;
while (e) {
struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)e;
if (nq->node_id == group->instance->node_id) { e = e->next; continue; }
if (!nq->paths) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "node has no paths"); e = e->next; continue; }
if (topo_group_should_send_to(nq, target) && topo_group_send_nodeinfo(group, nq, conn, topo_get_chain_rtt(nq)) < 0) return -1;
e = e->next;
}
return 0;
}
/* Обрабатывает REQUEST_TABLE: шлёт свой nodeinfo + полную таблицу + TABLE_COMPLETE. */
/* Повторный REQUEST_TABLE той же сессии не создаёт второй снимок. */
static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
if (!topo_group_peer_allowed(group, conn->peer_node_id)) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "table request deferred: member unknown peer=%016llx group=%016llx",
(unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id);
return;
}
if (topo_group_send_nodeinfo(group, group->local_node, conn, 0) < 0 || topo_group_send_full_table(group, conn) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table exchange incomplete group=%016llx peer=%016llx",
(unsigned long long)group->group_id, (unsigned long long)conn->peer_node_id);
return;
}
/* senders_list заполняется только через topo_group_new_conn (add + BGP);
здесь не добавляем, иначе new_conn (на JOIN_GROUP) упирается в дедуп и не шлёт TABLE_REQ обратно. */
if (topo_group_send_table_complete(group, conn) < 0) return;
struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id);
if (peer && peer->conn == conn && !peer->table_sent) { peer->table_sent = 1; topo_group_log_exchange(group, peer); }
if (peer && peer->accepted && !peer->table_sent) topo_tx_table(peer);
}
/* Обрабатывает JOIN_GROUP: для UTUN — добавить запросившего и инициировать BGP обратно. */

15
src/routing_layer/topo_group.h

@ -51,6 +51,7 @@ struct CONN_MGR;
struct TOPO_RECOVERY_CTX;
struct TOPO_GROUP_CONNECT;
struct TOPO_PEER_REQUEST;
struct topo_tx_item;
struct broadcast_ctx;
struct radio_ctx;
@ -106,6 +107,11 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g
* оба поколения и принимаются только в согласованной сессии. TABLE_COMPLETE
* отправляется после передачи транспорту всей таблицы отправителя.
* READY = живой транспорт + ACCEPT + table_sent + table_received.
* Исходящие команды одной сессии сохраняют FIFO-порядок. Снимок таблицы хранит
* только node_id; NODEINFO сериализуется при освобождении send_input_q. Удалённые
* узлы пропускаются, изменения после запроса идут за снимком. Повторный запрос
* не создаёт вторую таблицу. REINIT/leave/stop отменяют остаток и waiter; уже
* переданные транспорту пакеты защищены поколениями на принимающей стороне.
*
* Новое присоединение и REINIT создают новое поколение; REINIT сохраняет пути.
* Поколения защищают сообщения, а доступность существующего пути определяется
@ -173,6 +179,10 @@ struct TOPO_GROUP_CONN_ITEM {
uint8_t retained; // участие запрошено постоянным владельцем или достигло READY
uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения
uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch
uint8_t tx_failed;
struct topo_tx_item *tx_head, *tx_tail;
struct queue_waiter_handle tx_waiter;
void* tx_wake;
uint8_t table_received; // TABLE_COMPLETE текущей сессии принят
uint8_t table_sent; // ответ на запрос пира целиком передан транспорту
};
@ -341,9 +351,10 @@ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* han
int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* sender, const uint8_t* data, size_t len);
/**
* @brief Отправляет NODEINFO пакет одному conn (всегда local_node).
* @brief Планирует NODEINFO узла одному пиру с backpressure.
* Запись и RTT читаются при отправке; повторные изменения узла объединяются.
*/
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt);
int topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn);
/**
* @brief Отправляет WITHDRAW для node_id (вызывает broadcast_withdraw).

2
src/routing_layer/topo_node.c

@ -1207,7 +1207,7 @@ int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn);
se = se->next;
}
}

2
src/transport_layer/etcp_connections.c

@ -1906,7 +1906,7 @@ static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pk
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0);
if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn);
se = se->next;
}
}

4
tests/test_etcp_connect.c

@ -232,6 +232,10 @@ static void test11(void* arg) {
/* Shutdown обязан отменить ещё не отправленный LEAVE и освободить его waiter/handle. */
struct ll_queue* queue = a->send_input_q;
queue->callback_suspended = 1;
struct ll_entry* blocker = ll_alloc_lldgram(2);
if (!blocker) { fail("cannot allocate shutdown backpressure fixture"); return; }
blocker->dgram[0] = ETCP_ID_TOPO_ENTRY; blocker->dgram[1] = TOPO_SUBCMD_RESYNC; blocker->len = 2;
if (etcp_send(a, blocker) != 0) { queue_dgram_free(blocker); queue_entry_free(blocker); fail("cannot block shutdown queue"); return; }
struct TOPO_GROUP* group = topo_groups_get_default(g_a->topo_groups);
if (topo_group_join_peer(group, chat_a) < 0 || topo_group_leave_peer(group, chat_a) < 0 || !queue->waiter_head) {
fail("cannot prepare pending LEAVE shutdown fixture"); return;

16
tests/test_group_exchange.c

@ -14,6 +14,17 @@ static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) {
return (struct TOPO_GROUP_CONN_ITEM*)group->senders_list->head->data;
}
/* Deterministic transport consumer: consume one packet, then let the next producer run. */
static void drain(struct ETCP_CONN* conn) {
for (int i = 0; i < 128; i++) {
uasync_poll(conn->instance->ua, 0);
assert(queue_entry_count(conn->send_input_q) <= 1);
struct ll_entry* e = queue_data_get(conn->send_input_q);
if (e) { queue_dgram_free(e); queue_entry_free(e); }
queue_resume_callback(conn->send_input_q);
}
}
static void packet(struct ETCP_CONN* conn, const void* data, size_t size) {
struct ll_entry* e = ll_alloc_lldgram(size); assert(e);
memcpy(e->dgram, data, size); e->len = size;
@ -61,6 +72,7 @@ static void complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t
static void request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
receive(conn, group->group_id, TOPO_SUBCMD_REQUEST_TABLE, peer(group)->local_epoch, sizeof(struct TOPOMSG_HEADER));
drain(conn);
}
int main(void) {
@ -86,6 +98,7 @@ int main(void) {
struct ETCP_CONN* conn = node_conn_direct_get_conn(owner); assert(conn);
/* Детерминированная доставка управляющих пакетов через штатный BGP receiver.
* Event loop не запускаем: транспортный handshake проверяют интеграционные тесты. */
queue_set_callback(conn->send_input_q, NULL, NULL);
conn->links_up = 1;
assert(topo_group_new_conn(a, conn) == 0 && topo_group_new_conn(b, conn) == 0);
uint64_t a_id = peer(a)->local_epoch, b_id = peer(b)->local_epoch;
@ -125,7 +138,10 @@ int main(void) {
conn->send_input_q->size_limit = 0;
request(a, conn);
assert(!peer(a)->table_sent);
assert(peer(a)->tx_failed);
conn->send_input_q->size_limit = saved_limit;
etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT);
accept_join(a, conn);
complete(a, conn, peer(a)->local_epoch); request(a, conn);
assert(topo_group_peer_ready(a, node.node_id));

73
tests/test_group_recovery.c

@ -46,6 +46,8 @@ static void create(struct fixture* f) {
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;
assert(sc_init_ctx(&conn->crypto_ctx, &f->inst->my_keys) == SC_OK);
assert(sc_set_peer_public_key(&conn->crypto_ctx, node->public_key, SC_PEER_PUBKEY_BIN) == SC_OK);
queue_set_callback(conn->send_input_q, NULL, NULL);
etcp_conn_ready(conn);
}
@ -91,11 +93,24 @@ static void accept_join(struct fixture* f, int i) {
assert(peer(f, i)->accepted);
}
static void drain(struct fixture* f, int index) {
for (int i = 0; i < 128; i++) {
uasync_poll(f->ua, 0);
struct TOPO_GROUP_CONN_ITEM* p = peer(f, index);
if (!p || !p->conn) return;
assert(queue_entry_count(p->conn->send_input_q) <= 1);
struct ll_entry* e = queue_data_get(p->conn->send_input_q);
if (e) { queue_dgram_free(e); queue_entry_free(e); }
queue_resume_callback(p->conn->send_input_q);
}
}
static void ready(struct fixture* f, int i) {
assert(peer(f, i) && peer(f, i)->local_epoch);
accept_join(f, i);
table(f, i, TOPO_SUBCMD_REQUEST_TABLE, peer(f, i)->local_epoch);
table(f, i, TOPO_SUBCMD_TABLE_COMPLETE, peer(f, i)->local_epoch);
drain(f, i);
assert(topo_group_peer_ready(f->group, f->ids[i]));
}
@ -136,10 +151,11 @@ static void stalled_and_exhausted(void) {
topo_recovery_start(f.group); poll_events(&f);
up(&f, 0);
accept_join(&f, 0);
drain(&f, 0); /* Complete the preceding UP/ACCEPT callbacks before simulating a stalled exchange. */
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, peer(&f, 0)->local_epoch); poll_events(&f);
table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch); drain(&f, 0); 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);
@ -195,9 +211,64 @@ static void connect_deadline(void) {
destroy(&f);
}
static void sender_backpressure(void) {
struct fixture f; create(&f);
f.conns[0]->links_up = 1; f.conns[2]->links_up = 1;
assert(topo_group_new_conn(f.group, f.conns[0]) == 0); accept_join(&f, 0);
route(&f, 1, 2); route(&f, 2, 2);
table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch);
table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch); /* one snapshot */
for (int i = 0; i < 50; i++) {
assert(topo_group_send_nodeinfo(f.group, f.group->local_node, f.conns[0]) == 0);
uasync_poll(f.ua, 0);
}
assert(!peer(&f, 0)->table_sent && queue_entry_count(f.conns[0]->send_input_q) == 1);
struct TOPO_GROUP_NODE* removed = topo_node_find_by_id(f.group, f.ids[1]);
topo_group_remove_path(removed, f.conns[2]); topo_nodeq_remove_node(f.group, removed);
int local = 0, remote = 0, complete = 0;
for (int i = 0; i < 128; i++) {
uasync_poll(f.ua, 0);
assert(queue_entry_count(f.conns[0]->send_input_q) <= 1);
struct ll_entry* e = queue_data_get(f.conns[0]->send_input_q);
if (e && e->dgram[0] == ETCP_ID_TOPO_ENTRY) {
struct TOPOMSG_HEADER* h = (struct TOPOMSG_HEADER*)e->dgram;
if (h->subcmd == TOPO_SUBCMD_NODEINFO) {
uint64_t id = ((struct TOPOMSG_NODEINFO_PKT*)e->dgram)->node.node_id;
assert(id != f.ids[1]); /* removed during backpressure: never serialize a stale pointer */
if (id == f.inst->node_id) local++; else { assert(id == f.ids[2]); remote++; }
} else if (h->subcmd == TOPO_SUBCMD_TABLE_COMPLETE) {
assert(local == 1 && remote == 1); complete++;
}
}
if (e) { queue_dgram_free(e); queue_entry_free(e); }
queue_resume_callback(f.conns[0]->send_input_q);
}
assert(complete == 1 && local == 2 && remote == 1 && peer(&f, 0)->table_sent);
/* 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);
table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch);
uint64_t old = peer(&f, 0)->local_epoch;
etcp_fire_conn_status(f.conns[0], ETCP_CONN_STATUS_REINIT);
for (int i = 0; i < 32; i++) {
uasync_poll(f.ua, 0);
struct ll_entry* e = queue_data_get(f.conns[0]->send_input_q);
if (e && e->dgram[0] == ETCP_ID_TOPO_ENTRY) {
struct TOPOMSG_HEADER* h = (struct TOPOMSG_HEADER*)e->dgram;
assert(h->src_epoch != old && h->subcmd == TOPO_SUBCMD_JOIN_GROUP);
}
if (e) { queue_dgram_free(e); queue_entry_free(e); }
queue_resume_callback(f.conns[0]->send_input_q);
}
accept_join(&f, 0); table(&f, 0, TOPO_SUBCMD_REQUEST_TABLE, peer(&f, 0)->local_epoch);
poll_events(&f);
topo_groups_remove_group(f.inst->topo_groups, f.group->group_id); f.group = NULL;
poll_events(&f); destroy(&f);
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
utun_instance_set_tun_init_enabled(0);
sender_backpressure();
partial_and_external_routes();
stalled_and_exhausted();
cancel_and_group_stop();

Loading…
Cancel
Save