#include #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; } /* 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; conn->instance->api_bindings.callbacks[ETCP_ID_TOPO_ENTRY](conn, e); } static void nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* conn, struct TOPO_NODE* node, struct SC_MYKEYS* keys, uint64_t generation, uint64_t timestamp) { node->timestamp = timestamp; EVP_PKEY* key = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, keys->private_key, 32); size_t key_len = 32; assert(key && EVP_PKEY_get_raw_public_key(key, node->ed25519_public_key, &key_len) == 1); EVP_PKEY_free(key); uint8_t signed_data[TOPO_SIG_MSG_MAX_SIZE]; int length = topo_node_build_sig_msg(node, signed_data, sizeof(signed_data)); assert(length > 0); assert(sc_ed25519_sign(keys->private_key, signed_data, length, node->x25519_self_sig) == SC_OK); struct TOPOMSG_HEADER h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_NODEINFO, .group_id = group->group_id, .src_epoch = group->group_id + 1000, .dst_epoch = generation }; uint8_t wire[TOPO_NODE_WIRE_MAX_SIZE + sizeof(struct TOPOMSG_NODEINFO_PKT)]; memcpy(wire, &h, sizeof(h)); struct TOPO_GROUP_NODE nq = {0}; length = topo_node_serialize(node, &nq, group->group_id, TOPO_FLAG_SEND_SUBNETS, wire + sizeof(h), sizeof(wire) - sizeof(h), 0); assert(length > 0); packet(conn, wire, sizeof(h) + length); } static void receive(struct ETCP_CONN* conn, uint64_t group_id, uint8_t cmd, uint64_t local_epoch, size_t length) { struct TOPOMSG_JOIN_GROUP msg = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = cmd, .group_id = group_id, .src_epoch = group_id + 1000, .dst_epoch = local_epoch }, .group_type = TOPO_GROUP_TYPE_UTUN }; 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 accept_join(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { receive(conn, group->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); receive(conn, group->group_id, TOPO_SUBCMD_JOIN_ACCEPT, peer(group)->local_epoch, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(peer(group)->accepted); } static void complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t local_epoch) { receive(conn, group->group_id, TOPO_SUBCMD_TABLE_COMPLETE, local_epoch, sizeof(struct TOPOMSG_HEADER)); } 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) { 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 проверяют интеграционные тесты. */ 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; 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, a_id); assert(!peer(a)->table_received); /* no topology before ACCEPT */ accept_join(a, conn); accept_join(b, conn); 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)->local_epoch == a_id); receive(conn, a->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(peer(a)->local_epoch == 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)); nodeinfo(a, conn, &node, &keys, a_id, 1); assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 1); /* Отзыв альтернативы не должен уничтожать сохранившийся прямой путь. */ struct ETCP_CONN relay = { .instance = inst, .peer_node_id = 0x1234, .links_up = 1 }; uint64_t relay_hops[] = { node.node_id, relay.peer_node_id }; struct TOPO_GROUP_NODE* routed = topo_node_find_by_id(a, node.node_id); assert(topo_group_add_path(routed, &relay, relay_hops, 2, 10) == 0); struct TOPOMSG_WITHDRAW_PKT alternate_wd = { .node_id = node.node_id, .wd_source = relay.peer_node_id }; DEBUG_INFO(DEBUG_CATEGORY_BGP, "regression: withdraw one of two paths node=%016llx paths=%d", (unsigned long long)node.node_id, queue_entry_count(routed->paths)); assert(topo_group_process_withdraw(a, &relay, (const uint8_t*)&alternate_wd, sizeof(alternate_wd)) == 0); assert(topo_node_find_by_id(a, node.node_id) == routed); assert(queue_entry_count(routed->paths) == 1 && topo_group_find_conn_for_node(a, node.node_id) == conn); assert(topo_group_process_withdraw(a, &relay, (const uint8_t*)&alternate_wd, sizeof(alternate_wd)) == 0); assert(topo_node_find_by_id(a, node.node_id) == routed && queue_entry_count(routed->paths) == 1); etcp_fire_conn_status(conn, ETCP_CONN_STATUS_REINIT); assert(topo_group_find_conn_for_node(a, node.node_id) == conn); /* REINIT preserves known paths */ assert(peer(a)->local_epoch != a_id && peer(b)->local_epoch != 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); accept_join(a, conn); /* Ошибка отправки снимка не разрешает 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); 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)); nodeinfo(a, conn, &node, &keys, a_id, 2); assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 1); nodeinfo(a, conn, &node, &keys, peer(a)->local_epoch, 2); assert(topo_node_find_by_id(a, node.node_id)->last_timestamp == 2); struct TOPOMSG_WITHDRAW_PKT wd = { .h = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_WITHDRAW, .group_id = a->group_id, .src_epoch = a->group_id + 1000, .dst_epoch = a_id }, .node_id = node.node_id, .wd_source = node.node_id }; packet(conn, &wd, sizeof(wd)); assert(topo_node_find_by_id(a, node.node_id)); wd.h.dst_epoch = peer(a)->local_epoch; packet(conn, &wd, sizeof(wd)); assert(!topo_node_find_by_id(a, node.node_id)); a_id = peer(a)->local_epoch; 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)->local_epoch != a_id); receive(conn, a->group_id, TOPO_SUBCMD_JOIN_ACCEPT, a_id, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(!peer(a)->accepted); accept_join(a, conn); receive(conn, a->group_id, TOPO_SUBCMD_LEAVE_GROUP, a_id, sizeof(struct TOPOMSG_HEADER)); assert(peer(a)->accepted); /* old LEAVE cannot detach the new session */ complete(a, conn, a_id); request(a, conn); assert(!topo_group_peer_ready(a, node.node_id)); complete(a, conn, peer(a)->local_epoch); assert(topo_group_peer_ready(a, node.node_id) && !topo_group_peer_ready(b, node.node_id)); receive(conn, b->group_id, TOPO_SUBCMD_JOIN_REJECT, b_id, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(b->senders_list->head); receive(conn, b->group_id, TOPO_SUBCMD_JOIN_REJECT, peer(b)->local_epoch, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(!b->senders_list->head && !conn->close_requested); /* Cancellation while only JOIN has been received must not leave a half-open peer. */ receive(conn, b->group_id, TOPO_SUBCMD_JOIN_GROUP, 0, sizeof(struct TOPOMSG_JOIN_GROUP)); assert(peer(b)->peer_epoch && !peer(b)->accepted); receive(conn, b->group_id, TOPO_SUBCMD_LEAVE_GROUP, 0, sizeof(struct TOPOMSG_HEADER)); assert(!b->senders_list->head); 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; }