diff --git a/src/chat/chat_status.c b/src/chat/chat_status.c index d8335bff..eb51633a 100644 --- a/src/chat/chat_status.c +++ b/src/chat/chat_status.c @@ -218,7 +218,7 @@ static void collect_conn_metrics(struct UTUN_INSTANCE* inst, uint64_t peer_node_ "=== PEER: %04llX\u2192%04llX [%s] ===\n" "Role: %s\n" "Status: %s Links: %d/%d Initialized: %s MTU: %d\n" - "Reinit: %u Reset: %u Tx_state: %d Routing_exchange: %d\n" + "Reinit: %u Reset: %u Tx_state: %d\n" "\n--- ETCP ---\n" "RTT last/avg10: %.1f/%.1f ms Jitter: %.1f ms\n" "Bytes sent: %u Retrans: %u ACKs: %u\n" @@ -230,7 +230,7 @@ static void collect_conn_metrics(struct UTUN_INSTANCE* inst, uint64_t peer_node_ role_str, conn->links_up ? "UP" : "DOWN", links_up_count, link_count, conn->initialized ? "yes" : "no", conn->mtu, - conn->reinit_count, conn->reset_count, conn->tx_state, conn->routing_exchange_active, + conn->reinit_count, conn->reset_count, conn->tx_state, conn->rtt_last / 10.0f, conn->rtt_avg_10 / 10.0f, conn->jitter / 10.0f, conn->bytes_sent_total, conn->retransmissions_count, conn->ack_packets_count, conn->unacked_bytes, conn->max_inflight, diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index e5e43f7b..bdf39816 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -37,19 +37,47 @@ // Вспомогательные функции // ============================================================================ +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) { + 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; + } + return NULL; +} + +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); + return peer && peer->conn->links_up && peer->table_received && peer->table_sent; +} + +static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_CONN_ITEM* peer) { + 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)peer->exchange_id, peer->table_sent, peer->table_received, + topo_group_peer_ready(group, peer->conn->peer_node_id)); +} + /* Шлёт пиру запрос полной таблицы узлов группы (REQUEST_TABLE). */ static void topo_group_send_table_request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { if (!group || !conn) return; + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, conn->peer_node_id); + if (!peer || peer->conn != conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request without group peer"); return; } + if (!peer->exchange_id && (random_bytes((uint8_t*)&peer->exchange_id, sizeof(peer->exchange_id)) != 0 || !peer->exchange_id)) { + peer->exchange_id = 0; + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot generate group exchange id group=%016llx", (unsigned long long)group->group_id); + return; + } DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Sending table request to %s grp=%016llx", conn->log_name, (unsigned long long)group->group_id); struct TOPOMSG_TABLE_REQ* req = u_calloc(1, sizeof(struct TOPOMSG_TABLE_REQ)); - if (!req) return; + if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request allocation failed"); return; } req->cmd = ETCP_ID_TOPO_ENTRY; req->subcmd = TOPO_SUBCMD_REQUEST_TABLE; req->group_id = group->group_id; + req->exchange_id = peer->exchange_id; struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(req); return; } + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table request entry allocation failed"); u_free(req); return; } e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_TABLE_REQ); - if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } + if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send table request"); u_free(req); queue_entry_free(e); } } /* Запрашивает у пира членство в группе (JOIN_GROUP). */ @@ -82,32 +110,37 @@ static void topo_group_send_resync(struct ETCP_CONN* conn) { } /* Сигнализирует пиру об окончании начальной синхронизации таблицы. */ -static void topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || !conn) return; +static int topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { struct TOPOMSG_TABLE_REQ* req = u_calloc(1, sizeof(struct TOPOMSG_TABLE_REQ)); - if (!req) return; + if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table complete allocation failed"); return -1; } req->cmd = ETCP_ID_TOPO_ENTRY; req->subcmd = TOPO_SUBCMD_TABLE_COMPLETE; req->group_id = group->group_id; + req->exchange_id = exchange_id; struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(req); return; } + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "table complete entry allocation failed"); u_free(req); return -1; } e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_TABLE_REQ); - if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } + if (etcp_send(conn, e) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot send table complete"); u_free(req); queue_entry_free(e); return -1; + } + return 0; } 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 void 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 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, uint64_t exchange_id); static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn); static void topo_group_send_resync(struct ETCP_CONN* conn); static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn); /* Запоминает REQUEST_TABLE для группы, которой у нас ещё нет (CHAT-канал не загружен). */ -static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_id, uint64_t group_id) { +static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_id, uint64_t group_id, uint64_t exchange_id) { if (!g) return; for (int i = 0; i < g->pending_table_req_count; i++) - if (g->pending_table_reqs[i].node_id == node_id && g->pending_table_reqs[i].group_id == group_id) return; + if (g->pending_table_reqs[i].node_id == node_id && g->pending_table_reqs[i].group_id == group_id) { + g->pending_table_reqs[i].exchange_id = exchange_id; return; + } if (g->pending_table_req_count >= TOPO_MAX_PENDING_TABLE_REQS) { memmove(g->pending_table_reqs, g->pending_table_reqs + 1, (TOPO_MAX_PENDING_TABLE_REQS - 1) * sizeof(g->pending_table_reqs[0])); @@ -115,6 +148,7 @@ static void topo_group_remember_table_req(struct TOPO_GROUPS* g, uint64_t node_i } g->pending_table_reqs[g->pending_table_req_count].node_id = node_id; g->pending_table_reqs[g->pending_table_req_count].group_id = group_id; + g->pending_table_reqs[g->pending_table_req_count].exchange_id = exchange_id; g->pending_table_req_count++; DEBUG_INFO(DEBUG_CATEGORY_BGP, "remembered table request node=%016llx grp=%016llx (group not created yet)", (unsigned long long)node_id, (unsigned long long)group_id); @@ -133,7 +167,7 @@ static void topo_group_fulfill_table_reqs(struct TOPO_GROUPS* g, uint64_t group_ if (conn && conn->links_up) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "fulfilling remembered table request node=%016llx grp=%016llx", (unsigned long long)req->node_id, (unsigned long long)group_id); - topo_group_handle_request_table(group, conn); + topo_group_handle_request_table(group, conn, req->exchange_id); } } g->pending_table_req_count = w; @@ -259,6 +293,13 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry* queue_dgram_free(entry); queue_entry_free(entry); return; } + if ((subcmd == TOPO_SUBCMD_REQUEST_TABLE || subcmd == TOPO_SUBCMD_TABLE_COMPLETE) && + (entry->len != sizeof(struct TOPOMSG_TABLE_REQ) || !((struct TOPOMSG_TABLE_REQ*)data)->exchange_id)) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid table control peer=%016llx cmd=%u length=%u", + (unsigned long long)from_conn->peer_node_id, subcmd, entry->len); + queue_dgram_free(entry); queue_entry_free(entry); return; + } + uint64_t pkt_group_id = 0; if (subcmd == TOPO_SUBCMD_NODEINFO && entry->len >= sizeof(struct TOPOMSG_NODEINFO_PKT)) { pkt_group_id = ((struct TOPOMSG_NODEINFO_PKT*)data)->node.group_id; @@ -278,7 +319,8 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry* struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, pkt_group_id); if (!group) { if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) - topo_group_remember_table_req(instance->topo_groups, from_conn->peer_node_id, pkt_group_id); + topo_group_remember_table_req(instance->topo_groups, from_conn->peer_node_id, pkt_group_id, + ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id); DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "BGP recv %s from %s: group %016llx not found, dropping", group_subcmd_name(subcmd), from_conn->log_name, (unsigned long long)pkt_group_id); queue_dgram_free(entry); queue_entry_free(entry); return; } @@ -302,14 +344,25 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry* if (subcmd == TOPO_SUBCMD_NODEINFO) { nodeinfo_dump_log(data, entry->len); topo_group_process_nodeinfo(group, from_conn, data, entry->len); } else if (subcmd == TOPO_SUBCMD_WITHDRAW) topo_group_process_withdraw(group, from_conn, data, entry->len); - else if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) topo_group_handle_request_table(group, from_conn); + else if (subcmd == TOPO_SUBCMD_REQUEST_TABLE) + topo_group_handle_request_table(group, from_conn, ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id); else if (subcmd == TOPO_SUBCMD_JOIN_GROUP) topo_group_handle_join_group(group, from_conn); else if (subcmd == TOPO_SUBCMD_LEAVE_GROUP) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "peer left group: peer=%016llx group=%016llx", (unsigned long long)from_conn->peer_node_id, (unsigned long long)group->group_id); topo_group_remove_conn(group, from_conn, TOPO_REMOVE_REMOTE_LEAVE); } - else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE) { etcp_set_routing_exchange_state(from_conn, 3); DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP sync complete with %s: %d nodes grp=%016llx", from_conn->log_name, queue_entry_count(group->nodes), (unsigned long long)group->group_id); } + else if (subcmd == TOPO_SUBCMD_TABLE_COMPLETE) { + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, from_conn->peer_node_id); + uint64_t id = ((struct TOPOMSG_TABLE_REQ*)data)->exchange_id; + if (!peer || peer->conn != from_conn || peer->exchange_id != id) { + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ignore stale TABLE_COMPLETE group=%016llx peer=%016llx exchange=%016llx", + (unsigned long long)group->group_id, (unsigned long long)from_conn->peer_node_id, (unsigned long long)id); + } else if (!peer->table_received) { + peer->table_received = 1; + topo_group_log_exchange(group, peer); + } + } else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH) { if (entry->len >= sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)) { struct TOPOMSG_ERR_GROUP_MISMATCH* err = (struct TOPOMSG_ERR_GROUP_MISMATCH*)data; @@ -336,7 +389,9 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg for (struct ll_entry* e = groups->group_list->head; e; e = e->next) { struct TOPO_GROUP* g = (struct TOPO_GROUP*)e; for (struct ll_entry* se = g->senders_list->head; se; se = se->next) { - if (((struct TOPO_GROUP_CONN_ITEM*)se->data)->conn != conn) continue; + struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)se->data; + if (peer->conn != conn) continue; + peer->exchange_id = 0; peer->table_received = 0; peer->table_sent = 0; 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); topo_group_send_join_group(g,conn); topo_group_send_table_request(g,conn); @@ -1163,12 +1218,12 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send } /* Сериализует и шлёт NODEINFO узла конкретному пиру. */ -void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, struct ETCP_CONN* conn, uint16_t cumulative_rtt) { - if (!group || !node || !conn) return; +int topo_group_send_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; uint8_t* p = u_malloc(max_sz); - if (!p) return; + if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: allocation failed"); return -1; } p[0] = ETCP_ID_TOPO_ENTRY; p[1] = TOPO_SUBCMD_NODEINFO; @@ -1183,18 +1238,26 @@ void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node->paths ? queue_entry_count(node->paths) : -1, group->nodes ? queue_entry_count(group->nodes) : -1); u_free(p); - return; + return -1; } DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx ver=%d grp=%016llx to conn=%p name='%s'", (unsigned long long)node->node_id, sni->ver, (unsigned long long)group->group_id, (void*)conn, conn->log_name); uint8_t sflags = (group->group_type == TOPO_GROUP_TYPE_CHAT) ? 0 : TOPO_FLAG_SEND_SUBNETS; int ser_len = topo_node_serialize(sni, node, group->group_id, sflags, p + 2, max_sz - 2, cumulative_rtt); - if (ser_len < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: serialize failed for node %016llx", (unsigned long long)node->node_id); u_free(p); return; } + if (ser_len < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: serialize failed for node %016llx", (unsigned long long)node->node_id); + u_free(p); return -1; + } struct ll_entry* e = queue_entry_new(0); - if (!e) { u_free(p); return; } + if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: entry allocation failed"); u_free(p); return -1; } e->dgram = p; e->len = (size_t)ser_len + 2; - if (etcp_send(conn, e) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s", (unsigned long long)node->node_id, conn->log_name); u_free(p); queue_entry_free(e); } + if (etcp_send(conn, e) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: etcp_send FAILED for node %016llx to %s", + (unsigned long long)node->node_id, conn->log_name); + u_free(p); queue_entry_free(e); return -1; + } + return 0; } /* Добавляет conn в senders_list (дедуп), если его там ещё нет. */ @@ -1315,32 +1378,37 @@ int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* han } /* Шлёт пиру полную таблицу узлов группы (все узлы, кроме достижимых через него). */ -static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || !conn) return; +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)); + 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. */ -static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { +static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { 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; } - topo_group_send_nodeinfo(group, group->local_node, conn, 0); - topo_group_send_full_table(group, conn); + 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 exchange=%016llx", + (unsigned long long)group->group_id, (unsigned long long)conn->peer_node_id, (unsigned long long)exchange_id); + return; + } /* senders_list заполняется только через topo_group_new_conn (add + BGP); здесь не добавляем, иначе new_conn (на JOIN_GROUP) упирается в дедуп и не шлёт TABLE_REQ обратно. */ - topo_group_send_table_complete(group, conn); + if (topo_group_send_table_complete(group, conn, exchange_id) < 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); } } /* Обрабатывает JOIN_GROUP: для UTUN — добавить запросившего и инициировать BGP обратно. */ diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index b2cc63ce..c799c836 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -98,6 +98,12 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g * JOIN_READY из chat_sync.h — результат добавления участника, а TABLE_COMPLETE — * результат обмена топологией одной группы; ни один не означает готовности других групп. * Конфиг, auto-connect и recovery обязаны использовать общую проверку членства. + * + * REQUEST_TABLE содержит ненулевой exchange_id, сгенерированный запросившей + * стороной для текущего обмена (группа, пир). TABLE_COMPLETE возвращает тот же + * id только после успешной постановки всех NODEINFO в транспорт. Дубликаты + * запроса используют прежний id; новое присоединение и REINIT создают новый. + * Ответ с другим id не завершает текущий обмен. Формат изменён без совместимости. */ // Sub-команды @@ -140,6 +146,7 @@ struct TOPOMSG_TABLE_REQ { uint8_t cmd; uint8_t subcmd; uint64_t group_id; // идентификатор группы + uint64_t exchange_id; // ненулевой id запроса; TABLE_COMPLETE возвращает его без изменений } __attribute__((packed)); /** @@ -173,9 +180,11 @@ struct TOPOMSG_ERR_GROUP_MISMATCH { } __attribute__((packed)); struct TOPO_GROUP_CONN_ITEM { - struct ll_entry ll; struct ETCP_CONN* conn; struct NODE_CONN_DIRECT* handle; // владение BGP этим conn (NCD-handle) + uint64_t exchange_id; // текущий исходящий REQUEST_TABLE + uint8_t table_received; // TABLE_COMPLETE для exchange_id принят + uint8_t table_sent; // ответ на запрос пира целиком передан транспорту }; struct ETCP_LINK; @@ -194,7 +203,7 @@ struct TOPO_GROUP { uint64_t group_id; // уникальный идентификатор группы uint8_t group_type; // TOPO_GROUP_TYPE_* struct UTUN_INSTANCE* instance; - struct ll_queue* senders_list; // TOPO_GROUP_CONN_ITEM{ll_entry,ETCP_CONN*} — активные BGP-пиры, без хеш-индекса + struct ll_queue* senders_list; // ll_entry.data = TOPO_GROUP_CONN_ITEM, без хеш-индекса struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B) struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе uint8_t ed25519_public_key[SC_PUBKEY_SIZE]; @@ -215,6 +224,7 @@ struct TOPO_GROUP { struct topo_pending_table_req { uint64_t node_id; /* кто запросил таблицу */ uint64_t group_id; /* какую группу */ + uint64_t exchange_id; }; /** @@ -294,6 +304,13 @@ void topo_groups_remove_group(struct TOPO_GROUPS* g, uint64_t group_id); */ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn); +/* Готовность относится только к паре (group, peer): собственный снимок отправлен, + * снимок пира принят до TABLE_COMPLETE текущего REQUEST_TABLE. Транспортный UP + * этого не гарантирует. Повторное new_conn не сбрасывает готовность; REINIT + * начинает новый обмен с новым exchange_id. Старый TABLE_COMPLETE игнорируется. + * Это состояние обмена таблицами, а не добавление мембера (см. chat_join.h). */ +int topo_group_peer_ready(const struct TOPO_GROUP* group, uint64_t peer_id); + /** * @brief Удаляет conn из senders_list, очищает paths во всех nodes, отправляет withdraw если node unreachable. * @@ -332,7 +349,7 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send /** * @brief Отправляет NODEINFO пакет одному conn (всегда local_node). */ -void 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, uint16_t cumulative_rtt); /** * @brief Отправляет WITHDRAW для node_id (вызывает broadcast_withdraw). diff --git a/src/transport_layer/etcp.h b/src/transport_layer/etcp.h index ccef9a18..c8c60008 100644 --- a/src/transport_layer/etcp.h +++ b/src/transport_layer/etcp.h @@ -194,7 +194,6 @@ struct ETCP_CONN { // uint32_t total_packets_sent; // Total packets sent counter - Not used // Flags - uint8_t routing_exchange_active; // 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт, 3-завершён, 4-пропущен (нет BGP) uint8_t got_initial_pkt; // uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен) uint64_t reset_id; // локальная эпоха потока (меняется только при локальном reset) @@ -216,7 +215,6 @@ struct ETCP_CONN { // Unified callback chain with event mask (init/reinit/up/down/node_changed) struct etcp_cbk_entry* cbks; uint8_t reinit_pending; // 1 = reinit в процессе, ждём завершения - void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен) uint8_t fin_wait : 1; // close-pending: ожидаем подтверждение закрытия от remote void (*fin_wait_clear_cb)(struct ETCP_CONN* conn, void* arg); // вызывается при сбросе fin_wait @@ -329,4 +327,4 @@ void etcp_update_mtu(struct ETCP_CONN* etcp); } #endif -#endif // ETCP_H \ No newline at end of file +#endif // ETCP_H diff --git a/src/transport_layer/etcp_api.c b/src/transport_layer/etcp_api.c index 0b185a10..5858c668 100644 --- a/src/transport_layer/etcp_api.c +++ b/src/transport_layer/etcp_api.c @@ -155,12 +155,6 @@ void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event) { struct etcp_socket_cbk_entry* cbe = sock->instance->socket_cbks; while (cbe) { struct etcp_socket_cbk_entry* n = cbe->next; if (cbe->event_mask & event) cbe->fn(sock, event, cbe->arg); cbe = n; } } -void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) { - if (!conn || conn->state == 2) return; - conn->routing_exchange_active = new_state; - if (new_state >= 3 && conn->bgp_ready_cbk) - conn->bgp_ready_cbk(conn); -} int etcp_bind(struct UTUN_INSTANCE* inst, uint8_t id, etcp_recv_fn callback) { if (!inst || !callback || (unsigned)id >= ETCP_MAX_BINDINGS) return -1; diff --git a/src/transport_layer/etcp_api.h b/src/transport_layer/etcp_api.h index 8282de5c..5c21b2e5 100644 --- a/src/transport_layer/etcp_api.h +++ b/src/transport_layer/etcp_api.h @@ -265,7 +265,6 @@ void etcp_fire_link_status_cbk(struct ETCP_LINK* link, int old_state, int old_st // ---- Background connection initialization ---- #define ETCP_CONNECT_EARLY 1 #define ETCP_CONNECT_LATE 2 -#define ETCP_CONNECT_BGP_READY 4 // BGP-синхронизация завершена typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int type); @@ -276,11 +275,11 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t * @param node Узел из node registry (TOPO_GROUP_NODE) * @param cb Коллбэк завершения (сигнатура etcp_connect_callback_t) * @param arg Пользовательский аргумент - * @param flags Битовая маска этапов (ETCP_CONNECT_EARLY / LATE / BGP_READY) + * @param flags Битовая маска этапов (ETCP_CONNECT_EARLY / LATE) * @return 0 при успехе, -1 при ошибке * * @note Коллбэк вызывается на каждом запрошенном этапе с параметром - * type = ETCP_CONNECT_EARLY / ETCP_CONNECT_LATE / ETCP_CONNECT_BGP_READY + * type = ETCP_CONNECT_EARLY / ETCP_CONNECT_LATE * (или 0 при таймауте/неудаче). При уже готовом соединении EARLY+LATE * доставляются немедленно. */ @@ -397,16 +396,6 @@ void etcp_remove_socket_cbk(struct UTUN_INSTANCE* inst, etcp_socket_cbk_fn fn, v */ void etcp_socket_cbk_fire(struct ETCP_SOCKET* sock, int event); -/** - * @brief Установить состояние обмена маршрут-таблицами (BGP) - * - * @param conn ETCP соединение - * @param new_state 0-не активен, 1-надо инициировать (клиент), 2-обмен идёт, - * 3-завершён, 4-пропущен (нет BGP) - * - * @note При new_state >= 3 и наличии bgp_ready_cbk вызывается этот коллбэк - */ -void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state); /** * @brief Внутренняя функция etcp: Коллбэк для очередей output ll_queue normalizer diff --git a/src/transport_layer/etcp_api_doc.md b/src/transport_layer/etcp_api_doc.md index 797090f8..d576d71d 100644 --- a/src/transport_layer/etcp_api_doc.md +++ b/src/transport_layer/etcp_api_doc.md @@ -70,10 +70,10 @@ etcp_set_new_conn_cbk(instance, on_new_conn, my_data); ```c etcp_connect(instance, node, on_connect, arg, - ETCP_CONNECT_EARLY | ETCP_CONNECT_BGP_READY); + ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE); ``` -Флаги: `ETCP_CONNECT_EARLY` — раннее подключение, `ETCP_CONNECT_LATE` — отложенное, `ETCP_CONNECT_BGP_READY` — после BGP-синхронизации. +Флаги: `ETCP_CONNECT_EARLY` — раннее подключение, `ETCP_CONNECT_LATE` — отложенное. ## 3. API @@ -107,7 +107,6 @@ etcp_connect(instance, node, on_connect, arg, | `etcp_conn_add_down_cbk/remove_down_cbk` | Аналогично для down | | `etcp_set_new_conn_cbk(inst, fn, arg)` | Установить коллбэк на входящее соединение | | `etcp_add_new_conn_cbk/remove_new_conn_cbk` | Добавить/убрать из цепочки new_conn | -| `etcp_set_routing_exchange_state(conn, state)` | Установить состояние обмена маршрутами; при state≥3 вызывает `bgp_ready_cbk` | | `etcp_connect(inst, node, cb, arg, flags)` | Фоновое подключение к узлу | ### Командные ID (сырые, первый байт кодограммы) diff --git a/src/transport_layer/etcp_connect.c b/src/transport_layer/etcp_connect.c index 9d5d0f70..ac2c28f1 100644 --- a/src/transport_layer/etcp_connect.c +++ b/src/transport_layer/etcp_connect.c @@ -138,13 +138,6 @@ static void connect_create_links_v6(struct ETCP_CONNECT* ctx, struct TOPO_GROUP_ addr_count, sock_count, link_count, addr_count == 0 ? " (v6_addrs EMPTY)" : "", (unsigned long long)node->node_id); } -static void connect_bgp_ready_cb(struct ETCP_CONN* conn) { - struct ETCP_CONNECT* ctx = connect_find(conn->instance, conn->peer_node_id); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "[etcp_connect] bgp_ready_cb: node=%016llx ctx=%p ctx_done=%d ra=%d", - (unsigned long long)conn->peer_node_id, (void*)ctx, ctx ? ctx->done : -1, conn->routing_exchange_active); - if (!ctx || ctx->done) return; - connect_deliver(ctx, ETCP_CONNECT_BGP_READY, 1); -} static void connect_initial_timeout_cb(void* arg) { struct ETCP_CONNECT* ctx = (struct ETCP_CONNECT*)arg; @@ -229,7 +222,9 @@ static void connect_init_cb(struct ETCP_CONN* conn, int event, void* arg) { (voi int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node, etcp_connect_callback_t cb, void* arg, uint8_t flags) { - if (!inst || !node || !cb) return -1; + if (!inst || !node || !cb || !flags || (flags & ~(ETCP_CONNECT_EARLY | ETCP_CONNECT_LATE))) { + DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "invalid connect arguments/flags=%u", flags); return -1; + } struct TOPO_NODE* ni = topo_node_registry_find(inst->topo_groups, node->node_id); if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node not in registry"); return -1; } @@ -294,7 +289,6 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node, if (!cn) { u_free(ctx); return -1; } cn->cb = cb; cn->arg = arg; cn->flags = flags; ctx->cb_list = cn; etcp_conn_add_cbk(conn, connect_init_cb, ctx, ETCP_CBK_EVENT_INIT); - conn->bgp_ready_cbk = connect_bgp_ready_cb; ctx->initial_timer = uasync_set_timeout(inst->ua, (int)inst->etcp_connect_timeout_tb, ctx, connect_initial_timeout_cb, "etcp_connect_init"); ctx->next = inst->pending_connects; inst->pending_connects = ctx; @@ -318,25 +312,13 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node, DEBUG_INFO(DEBUG_CATEGORY_ETCP_CONNECT, "[etcp_connect] node 0x%016llx: existing ready conn, delivering EARLY+LATE", (unsigned long long)node_id); struct etcp_connect_cb_node* n = ctx->cb_list; - struct etcp_connect_cb_node* keep = NULL; while (n) { struct etcp_connect_cb_node* next = n->next; - if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY); - if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE); - if (n->flags & ETCP_CONNECT_BGP_READY) { - if (conn->routing_exchange_active >= 3) n->cb(n->arg, ctx->conn, ETCP_CONNECT_BGP_READY); - else { n->next = keep; keep = n; n = next; continue; } - } - u_free(n); - n = next; - } - ctx->cb_list = keep; - if (keep) { - conn->bgp_ready_cbk = connect_bgp_ready_cb; - ctx->next = inst->pending_connects; inst->pending_connects = ctx; - } else { - u_free(ctx); + if (n->flags & ETCP_CONNECT_EARLY) n->cb(n->arg, ctx->conn, ETCP_CONNECT_EARLY); + if (n->flags & ETCP_CONNECT_LATE) n->cb(n->arg, ctx->conn, ETCP_CONNECT_LATE); + u_free(n); n = next; } + u_free(ctx); return 0; } @@ -379,7 +361,6 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct TOPO_GROUP_NODE* node, } etcp_conn_add_cbk(conn, connect_init_cb, ctx, ETCP_CBK_EVENT_INIT); - conn->bgp_ready_cbk = connect_bgp_ready_cb; connect_create_links_v4(ctx, node); connect_create_links_v6(ctx, node); @@ -415,7 +396,7 @@ void etcp_connect_cancel_for_conn(struct UTUN_INSTANCE* inst, struct ETCP_CONN* } } -/* Полная отмена pending-коннекта по node_id: снимает connect_init_cb/bgp_ready_cbk +/* Полная отмена pending-коннекта по node_id: снимает connect_init_cb * с conn (conn остаётся жив, например, переиспользуется ncd), отменяет таймеры, * вынимает ctx из pending_connects и освобождает. Безопасно вызывать, если ctx нет. */ void etcp_connect_cancel_by_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { @@ -427,7 +408,6 @@ void etcp_connect_cancel_by_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { *pp = ctx->next; if (ctx->conn) { etcp_conn_remove_cbk(ctx->conn, connect_init_cb, ctx); - ctx->conn->bgp_ready_cbk = NULL; } if (ctx->initial_timer) { uasync_cancel_timeout(inst->ua, ctx->initial_timer); ctx->initial_timer = NULL; } if (ctx->settle_timer) { uasync_cancel_timeout(inst->ua, ctx->settle_timer); ctx->settle_timer = NULL; } diff --git a/src/transport_layer/etcp_connect.h b/src/transport_layer/etcp_connect.h index c6fe47ba..9a6b586c 100644 --- a/src/transport_layer/etcp_connect.h +++ b/src/transport_layer/etcp_connect.h @@ -25,7 +25,7 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t * @param cb коллбэк (etcp_connect_callback_t) * @param arg пользовательский аргумент, передаваемый в cb * @param flags битовая маска интересующих фаз: ETCP_CONNECT_EARLY(1) | - * ETCP_CONNECT_LATE(2) | ETCP_CONNECT_BGP_READY(4) + * ETCP_CONNECT_LATE(2) * @return 0 при успехе, -1 при ошибке (до вызова коллбэка) * * @@ -78,11 +78,7 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t * │ * └─ conn->state == 1 (ready, уже проинициализирован): * → connect_deliver(ctx, ETCP_CONNECT_LATE) сразу - * → если routing_exchange_active >= 3: - * connect_deliver(ctx, ETCP_CONNECT_BGP_READY) - * → если в cb_list остались неудовлетворённые коллбэки - * (ждут BGP_READY, который ещё не готов): - * ctx остаётся в pending_connects с conn->bgp_ready_cbk + * → контекст освобождается; готовность топологии принадлежит группе * * 4. !conn && !ctx — новое подключение (основной путь) * → etcp_connection_create(inst, NULL) @@ -118,10 +114,6 @@ typedef void (*etcp_connect_callback_t)(void* arg, struct ETCP_CONN* conn, int t * ├─ Убирает connect_init_cb из conn->cbks * └─ connect_deliver(ctx, ETCP_CONNECT_LATE) → ctx освобождается * - * connect_bgp_ready_cb (вызывается при routing_exchange_active >= 3): - * └─ connect_deliver(ctx, ETCP_CONNECT_BGP_READY) - * (может произойти в любой момент после установки соединения) - * * connect_initial_timeout_cb (initial_timer): * ├─ ctx->done = 1 * ├─ Закрывает TCP-линк diff --git a/src/transport_layer/etcp_connect_doc.md b/src/transport_layer/etcp_connect_doc.md index 3a1518a4..646eeecc 100644 --- a/src/transport_layer/etcp_connect_doc.md +++ b/src/transport_layer/etcp_connect_doc.md @@ -1,7 +1,7 @@ # ETCP Connect (etcp_connect) ## 1. Назначение -Управление исходящими ETCP-соединениями. Предоставляет асинхронный API `etcp_connect()` для установки соединения с удалённым узлом: создаёт ETCP_CONN, настраивает криптографию, поднимает UDP-линки и опционально TCP/STCP-транспорт, отслеживает прогресс установки с таймерами. Доставляет коллбэки о фазах готовности (EARLY, LATE, BGP_READY). +Управление исходящими ETCP-соединениями. Предоставляет асинхронный API `etcp_connect()` для установки соединения с удалённым узлом: создаёт ETCP_CONN, настраивает криптографию, поднимает UDP-линки и опционально TCP/STCP-транспорт, отслеживает прогресс установки с таймерами. Доставляет коллбэки о фазах готовности (EARLY, LATE). ## 2. Как пользоваться @@ -21,8 +21,6 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) { // первый линк готов (можно начинать обмен) if (type & ETCP_CONNECT_LATE) // settle-фаза завершена (все линки проверены, неудачные удалены) - if (type & ETCP_CONNECT_BGP_READY) - // BGP-синхронизация завершена (routing_exchange_active >= 3) } ``` @@ -30,8 +28,7 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) { 1. `etcp_connect()` — создаёт ETCP_CONN, crypto_ctx, UDP-линки и TCP/STCP 2. `connect_ready_cb()` — первый линк стал ready → EARLY-коллбэк → запуск settle-таймера 3. `connect_settle_timeout_cb()` — settle (min_rtt × 8, clamp 500..2000ms) → удаление неудачных линков → LATE-коллбэк -4. `connect_bgp_ready_cb()` — BGP завершён → BGP_READY-коллбэк -5. При таймауте (`connect_initial_timeout_cb`) — type=0 (conn=NULL), все ресурсы освобождаются +4. При таймауте (`connect_initial_timeout_cb`) — type=0 (conn=NULL), все ресурсы освобождаются **Нюансы:** - Если соединение уже установлено — коллбэки вызываются немедленно @@ -57,7 +54,6 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) { |---|---| | `ETCP_CONNECT_EARLY` (1) | Первый линк готов | | `ETCP_CONNECT_LATE` (2) | Settle-фаза завершена | -| `ETCP_CONNECT_BGP_READY` (4) | BGP-синхронизация завершена | ### Функции | Функция | Описание | @@ -69,3 +65,7 @@ void my_connect_cb(void* arg, struct ETCP_CONN* conn, int type) { 1. **Таймаут установки** (`connect_initial_timeout_cb`): stcp_link_close → etcp_connection_close → connect_cancel → type=0 2. **Settle-таймаут** (`connect_settle_timeout_cb`): удаление неудачных линков/stcp → LATE → connect_cancel 3. **Штатное закрытие** (`utun_instance_destroy`): фаза 1 (detach) → фаза 2 (deferred: `etcp_connect_cancel_for_conn`) + +Готовность BGP проверяется для конкретной группы через `topo_group_peer_ready(group, peer_id)`. +Общий транспорт не имеет состояния завершения группового обмена. Для владения +соединениями используйте NCD; создание транспорта само по себе не вступает в группу. diff --git a/src/transport_layer/etcp_dump.c b/src/transport_layer/etcp_dump.c index feb659bc..e190d8ff 100644 --- a/src/transport_layer/etcp_dump.c +++ b/src/transport_layer/etcp_dump.c @@ -56,9 +56,9 @@ void etcp_dump_conn_state(struct ETCP_CONN* conn) { DHDR("=== ETCP CONN [%s] ===", conn->log_name); - DLOG("GENERAL: peer=0x%016llx state=%d init=%d l_up=%d tx=%d mtu=%d rout_ex=%d", + DLOG("GENERAL: peer=0x%016llx state=%d init=%d l_up=%d tx=%d mtu=%d", (unsigned long long)conn->peer_node_id, conn->state, conn->initialized, conn->links_up, - conn->tx_state, conn->mtu, conn->routing_exchange_active); + conn->tx_state, conn->mtu); DLOG("IDS: next_tx=%u last_rx=%u last_del=%u rx_ack_till=%u got_init=%d reset_done=%d", conn->next_tx_id, conn->last_rx_id, conn->last_delivered_id, conn->rx_ack_till, diff --git a/tests/Makefile.am b/tests/Makefile.am index b6b8d2fc..00860404 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -65,6 +65,7 @@ check_PROGRAMS = \ test_etcp_connect \ test_node_conn_direct \ test_group_ownership \ + test_group_exchange \ test_node_snapshot \ test_ncd_config \ test_db_sync \ @@ -394,6 +395,10 @@ test_group_ownership_SOURCES = test_group_ownership.c test_group_ownership_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_group_ownership_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_group_exchange_SOURCES = test_group_exchange.c +test_group_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_group_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_node_snapshot_SOURCES = test_node_snapshot.c test_node_snapshot_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_node_snapshot_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_router_reconnect.c b/tests/test_etcp_router_reconnect.c index 753a6622..c1233693 100644 --- a/tests/test_etcp_router_reconnect.c +++ b/tests/test_etcp_router_reconnect.c @@ -140,7 +140,8 @@ static int send_one_pkt(void) { static int transit_routes_ready(void) { struct TOPO_GROUP* group = g_b ? topo_groups_get_default(g_b->topo_groups) : NULL; - return group && topo_group_find_conn_for_node(group, node_a) && topo_group_find_conn_for_node(group, node_c); + return group && topo_group_peer_ready(group, node_a) && topo_group_peer_ready(group, node_c) && + topo_group_find_conn_for_node(group, node_a) && topo_group_find_conn_for_node(group, node_c); } static const char* state_name(int s) { @@ -190,6 +191,7 @@ static void state_step(void) { { struct ETCP_CONN* direct = instance_find_conn(g_a, node_c); if (direct && direct->links_up && direct->initialized && + topo_group_peer_ready(topo_groups_get_default(g_a->topo_groups), node_c) && topo_group_find_conn_for_node(topo_groups_get_default(g_a->topo_groups), node_c) == direct) { fprintf(stderr,"RECOVERY: direct a-c path established without test reconnect\n"); g_phase_sent = 0; g_state = ST_PHASE3_SEND; diff --git a/tests/test_group_exchange.c b/tests/test_group_exchange.c new file mode 100644 index 00000000..3299a9cf --- /dev/null +++ b/tests/test_group_exchange.c @@ -0,0 +1,108 @@ +#include +#include +#include "utun_instance.h" +#include "routing_layer/topo_group.h" +#include "transport_layer/etcp.h" +#include "transport_layer/etcp_api.h" +#include "transport_layer/node_conn_direct.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" + +static struct TOPO_GROUP_CONN_ITEM* peer(struct TOPO_GROUP* group) { + assert(group->senders_list->head && !group->senders_list->head->next); + return (struct TOPO_GROUP_CONN_ITEM*)group->senders_list->head->data; +} + +static void receive(struct ETCP_CONN* conn, uint64_t group_id, uint8_t cmd, uint64_t exchange_id, size_t length) { + struct TOPOMSG_TABLE_REQ msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = cmd, + .group_id = group_id, .exchange_id = exchange_id }; + assert(length <= sizeof(msg)); + struct ll_entry* packet = ll_alloc_lldgram(sizeof(msg)); assert(packet); + memcpy(packet->dgram, &msg, sizeof(msg)); packet->len = length; + conn->instance->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](conn, packet); +} + +static void complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t exchange_id) { + receive(conn, group->group_id, TOPO_SUBCMD_TABLE_COMPLETE, exchange_id, sizeof(struct TOPOMSG_TABLE_REQ)); +} + +static void request(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + receive(conn, group->group_id, TOPO_SUBCMD_REQUEST_TABLE, 123, sizeof(struct TOPOMSG_TABLE_REQ)); +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO); + utun_instance_set_tun_init_enabled(0); + struct UASYNC* ua = uasync_create(); assert(ua); + struct UTUN_INSTANCE* inst = utun_instance_create_from_str(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(inst && utun_instance_init(inst) == 0); + struct TOPO_GROUP* a = topo_groups_create_group(inst->topo_groups, 42, TOPO_GROUP_TYPE_UTUN, NULL); + struct TOPO_GROUP* b = topo_groups_create_group(inst->topo_groups, 43, TOPO_GROUP_TYPE_UTUN, NULL); + assert(a && b); + struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK); + struct TOPO_ADDR4 addr = { .addr = {127, 0, 0, 1}, .port = 9, .protocol = TOPO_PROTO_UDP }; + struct TOPO_NODE node = { .v4_addrs = &addr }; + memcpy(node.public_key, keys.public_key, SC_PUBKEY_SIZE); + node.node_id = sc_derive_node_id_from_pubkey(node.public_key); + struct NODE_CONN_DIRECT* owner = NULL; + assert(node_conn_direct_open_node(inst, node.node_id, NULL, NULL, &owner, &node, NULL) == NCD_NEW); + struct ETCP_CONN* conn = node_conn_direct_get_conn(owner); assert(conn); + /* Детерминированная доставка управляющих пакетов через штатный BGP receiver. + * Event loop не запускаем: транспортный handshake проверяют интеграционные тесты. */ + 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)->exchange_id, b_id = peer(b)->exchange_id; + assert(a_id && b_id && a_id != b_id); + assert(!topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + complete(a, conn, b_id); + receive(conn, a->group_id, TOPO_SUBCMD_TABLE_COMPLETE, a_id, 2); + assert(!peer(a)->table_received); + complete(a, conn, a_id); + assert(peer(a)->table_received && !topo_group_peer_ready(a, node.node_id)); + request(a, conn); + assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + assert(topo_group_new_conn(a, conn) == 0 && peer(a)->exchange_id == a_id); + receive(conn, a->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); + assert(peer(a)->exchange_id == a_id); + assert(queue_entry_count(a->senders_list) == 1 && topo_group_peer_ready(a, node.node_id)); + /* Другой порядок доставки: сначала наш снимок, потом подтверждение чужого. */ + request(b, conn); + assert(!topo_group_peer_ready(b, node.node_id)); + complete(b, conn, b_id); + assert(topo_group_peer_ready(b, node.node_id)); + + etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); + assert(peer(a)->exchange_id != a_id && peer(b)->exchange_id != b_id); + assert(!topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + complete(a, conn, a_id); complete(b, conn, b_id); + assert(!peer(a)->table_received && !peer(b)->table_received); + /* Ошибка отправки снимка не разрешает TABLE_COMPLETE и READY. */ + int saved_limit = conn->send_input_q->size_limit; + conn->send_input_q->size_limit = 0; + request(a, conn); + assert(!peer(a)->table_sent); + conn->send_input_q->size_limit = saved_limit; + complete(a, conn, peer(a)->exchange_id); request(a, conn); + assert(topo_group_peer_ready(a, node.node_id)); + + a_id = peer(a)->exchange_id; + topo_group_remove_conn(a, conn, TOPO_REMOVE_LOCAL_LEAVE); + complete(a, conn, a_id); + assert(!topo_group_peer_ready(a, node.node_id)); + assert(topo_group_new_conn(a, conn) == 0 && peer(a)->exchange_id != a_id); + complete(a, conn, a_id); request(a, conn); + assert(!topo_group_peer_ready(a, node.node_id)); + complete(a, conn, peer(a)->exchange_id); + assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); + conn->links_up = 0; + assert(!topo_group_peer_ready(a, node.node_id)); + node_conn_direct_close(owner); + inst->running = 0; utun_instance_destroy(inst); + uasync_poll(ua, 0); uasync_destroy(ua, 0); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: isolation, duplicate attach, stale replies, send failure and REINIT passed"); + return 0; +}