From 0756a6e020c1271abfe60d25c725849da93b8df5 Mon Sep 17 00:00:00 2001 From: evgeny Date: Mon, 28 Sep 2026 02:45:00 +0300 Subject: [PATCH] Serialize group topology lazily through a cancellable backpressure sender --- src/routing_layer/topo_group.c | 269 +++++++++++++++++++++---- src/routing_layer/topo_group.h | 15 +- src/routing_layer/topo_node.c | 2 +- src/transport_layer/etcp_connections.c | 2 +- tests/test_etcp_connect.c | 4 + tests/test_group_exchange.c | 16 ++ tests/test_group_recovery.c | 73 ++++++- 7 files changed, 336 insertions(+), 45 deletions(-) diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index f11d018c..12fa4bd4 100644 --- a/src/routing_layer/topo_group.c +++ b/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 обратно. */ diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 740bb8aa..70aa12ce 100644 --- a/src/routing_layer/topo_group.h +++ b/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). diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index 46455591..f5e3e0e3 100644 --- a/src/routing_layer/topo_node.c +++ b/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; } } diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index f2b18b14..ba2b1599 100644 --- a/src/transport_layer/etcp_connections.c +++ b/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; } } diff --git a/tests/test_etcp_connect.c b/tests/test_etcp_connect.c index 131e9e9d..786735e3 100644 --- a/tests/test_etcp_connect.c +++ b/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; diff --git a/tests/test_group_exchange.c b/tests/test_group_exchange.c index 72901109..42f8ea38 100644 --- a/tests/test_group_exchange.c +++ b/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)); diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c index 7ccc3df1..e427f7da 100644 --- a/tests/test_group_recovery.c +++ b/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();