Browse Source

Route UTUN configuration and group ownership through NCD

proxy
evgeny 4 days ago
parent
commit
bce17d9cfc
  1. 24
      doc/node_snapshot.md
  2. 2
      src/chat/member_sync.c
  3. 156
      src/routing_layer/topo_group.c
  4. 12
      src/routing_layer/topo_group.h
  5. 4
      src/routing_layer/topo_node.c
  6. 7
      src/routing_layer/topo_node.h
  7. 2
      src/routing_layer/topo_node_sqlite.c
  8. 236
      src/transport_layer/etcp_connections.c
  9. 6
      src/transport_layer/etcp_connections.h
  10. 44
      src/transport_layer/node_conn_direct.c
  11. 64
      src/utun_instance.c
  12. 9
      tests/Makefile.am
  13. 31
      tests/test_bgp_route_exchange.c
  14. 27
      tests/test_bgp_triangle.c
  15. 65
      tests/test_connection_loss.h
  16. 41
      tests/test_etcp_connect.c
  17. 18
      tests/test_etcp_router.c
  18. 48
      tests/test_etcp_router_reconnect.c
  19. 27
      tests/test_group_ownership.c
  20. 96
      tests/test_ncd_config.c
  21. 7
      tests/test_reality_bgp.c

24
doc/node_snapshot.md

@ -34,3 +34,27 @@ Timestamp 0 обозначает bootstrap-сведения без подпис
`node_conn_direct_open_node` сверяет переданную запись с реестром/БД. Более
свежая подписанная запись принимается целиком, старая уступает уже известной.
Bootstrap-адреса не перекрывают подписанную запись даже при DOWN.
Конфигурационные подключения UTUN используют один адаптер
`init_client_connections` для старта и reload. Все адреса, TCP и REALITY
передаются NCD; адаптер не создаёт ETCP-линки. Новый набор handles открывается
до освобождения старого; при ошибке старые handles сохраняются. Конфигурация
удерживает один handle на клиента, независимо от количества его адресов.
Включение конфигурационного пира в UTUN-группу выполняет callback NCD UP.
Транспортный UP и наличие callback старого `etcp_connect` больше не включают
соединение в UTUN-группу автоматически. Входящий JOIN остаётся явным запросом
участия. При пересоздании UTUN-группы учитываются её конфигурационные handles.
Удаление клиента отправляет `LEAVE_GROUP` и освобождает удержание UTUN-группы
с обеих сторон. Другие владельцы NCD сохраняют транспорт. Поздние NODEINFO и
WITHDRAW от отсоединённого UTUN-пира игнорируются, чтобы не вернуть его маршруты.
LEAVE ждёт backpressure через queue waiter и временно удерживает собственный
NCD handle. Закрытие транспорта отменяет ожидание; новое участие в группе
делает ожидающий LEAVE устаревшим, и он не отправляется.
Причина удаления пира передаётся явно: только потеря транспорта запускает
recovery каскадных маршрутов. Локальный/удалённый LEAVE и удаление битой
member-записи не должны восстанавливать намеренно закрытое участие.
REALITY использует полное SNI до 255 символов; размер поля в записи узла
согласован с транспортом. Если локального TCP listen-сокета нет, NCD может
создать исходящий TCP-линк с выбором локального адреса операционной системой.

2
src/chat/member_sync.c

@ -855,7 +855,7 @@ int member_sync_verify_and_purge(struct UTUN_INSTANCE* inst, const char* ch_id)
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid);
if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) {
struct ETCP_CONN* c = instance_find_conn(inst, nid);
if (c) topo_group_remove_conn(g, c);
if (c) topo_group_remove_conn(g, c, TOPO_REMOVE_MEMBER_INVALID);
}
}
purged++;

156
src/routing_layer/topo_group.c

@ -192,6 +192,7 @@ static const char* group_subcmd_name(uint8_t subcmd) {
case TOPO_SUBCMD_ERR_GROUP_MISMATCH: return "ERR_GROUP_MISMATCH";
case TOPO_SUBCMD_JOIN_GROUP: return "JOIN_GROUP";
case TOPO_SUBCMD_RESYNC: return "RESYNC";
case TOPO_SUBCMD_LEAVE_GROUP: return "LEAVE_GROUP";
default: return "?";
}
}
@ -225,6 +226,12 @@ static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t nod
// Приём пакетов
// ============================================================================
static struct NODE_CONN_DIRECT* topo_config_handle(struct UTUN_INSTANCE* instance, uint64_t node_id) {
for (struct CONFIG_CONN_HANDLE* owner = instance->config_conn_handles; owner; owner = owner->next)
if (owner->node_id == node_id) return owner->handle;
return NULL;
}
static int topo_group_has_sender(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
for (struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; e; e = e->next)
if (((struct TOPO_GROUP_CONN_ITEM*)e->data)->conn == conn) return 1;
@ -253,7 +260,7 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
}
uint64_t pkt_group_id = 0;
if (subcmd == TOPO_SUBCMD_NODEINFO && entry->len >= 3) {
if (subcmd == TOPO_SUBCMD_NODEINFO && entry->len >= sizeof(struct TOPOMSG_NODEINFO_PKT)) {
pkt_group_id = ((struct TOPOMSG_NODEINFO_PKT*)data)->node.group_id;
} else if (subcmd == TOPO_SUBCMD_WITHDRAW && entry->len >= sizeof(struct TOPOMSG_WITHDRAW_PKT)) {
pkt_group_id = ((struct TOPOMSG_WITHDRAW_PKT*)data)->group_id;
@ -263,7 +270,8 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
pkt_group_id = ((struct TOPOMSG_TABLE_REQ*)data)->group_id;
} else if (subcmd == TOPO_SUBCMD_ERR_GROUP_MISMATCH && entry->len >= sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)) {
pkt_group_id = ((struct TOPOMSG_ERR_GROUP_MISMATCH*)data)->group_id;
} else if (subcmd == TOPO_SUBCMD_JOIN_GROUP && entry->len >= sizeof(struct TOPOMSG_JOIN_GROUP)) {
} else if ((subcmd == TOPO_SUBCMD_JOIN_GROUP || subcmd == TOPO_SUBCMD_LEAVE_GROUP) &&
entry->len == sizeof(struct TOPOMSG_JOIN_GROUP)) {
pkt_group_id = ((struct TOPOMSG_JOIN_GROUP*)data)->group_id;
}
@ -284,10 +292,23 @@ 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_LEAVE_GROUP || (group->group_type == TOPO_GROUP_TYPE_UTUN &&
(subcmd == TOPO_SUBCMD_NODEINFO || subcmd == TOPO_SUBCMD_WITHDRAW))) &&
!topo_group_has_sender(group, from_conn)) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "ignore %s from detached peer=%016llx group=%016llx",
group_subcmd_name(subcmd), (unsigned long long)from_conn->peer_node_id, (unsigned long long)pkt_group_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
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_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_ERR_GROUP_MISMATCH) {
if (entry->len >= sizeof(struct TOPOMSG_ERR_GROUP_MISMATCH)) {
@ -325,24 +346,9 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg
return;
}
if (status == ETCP_CONN_STATUS_UP && conn->peer_node_id) {
int is_client = conn->instance->config && config_peer_in_clients(conn->instance->config, conn->crypto_ctx.peer_public_key);
int explicit_conn = conn->bgp_ready_cbk != NULL;
/* пассивная сторона: просим пира заново анонсировать свои группы (после флэпа/переподключения) */
if (!is_client && !explicit_conn)
topo_group_send_resync(conn);
/* non-CHAT (UTUN): узлы из конфига (clients) + явные подключения (etcp_connect). */
struct ll_entry* fe = groups->group_list->head;
while (fe) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)fe; fe = fe->next;
if (g->group_type != TOPO_GROUP_TYPE_CHAT) {
if (is_client || explicit_conn) {
topo_group_new_conn(g, conn); /* добавить + инициировать BGP */
topo_group_send_join_group(g, conn); /* запросить членство у пира */
}
}
}
/* UTUN присоединяет пир через UP принадлежащего ему NCD handle.
* Сам по себе транспортный UP не означает участия в UTUN-группе. */
if (!topo_config_handle(conn->instance, conn->peer_node_id)) topo_group_send_resync(conn);
/* CHAT: только члены каналов (peers_* в БД). */
sqlite3* db = conn->instance ? conn->instance->topo_sqlite_db : NULL;
uint64_t* chs = NULL; int chn = 0;
@ -357,7 +363,7 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg
} else if (status == ETCP_CONN_STATUS_DOWN || status == ETCP_CONN_STATUS_DELETE) {
struct ll_entry* fe = groups->group_list->head;
while (fe) { struct TOPO_GROUP* g = (struct TOPO_GROUP*)fe; fe = fe->next;
topo_group_remove_conn(g, conn);
topo_group_remove_conn(g, conn, TOPO_REMOVE_TRANSPORT_DOWN);
}
}
}
@ -636,7 +642,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
for (uint32_t slot = 0; slot < connections->hash_size; slot++) {
struct ll_entry* entry = connections->hash_table[slot];
while (entry) {
struct conn_queue_entry* cqe = (struct conn_queue_entry*)entry;
struct conn_queue_entry* cqe = (struct conn_queue_entry*)entry->data;
if (cqe->conn && cqe->peer_node_id != 0) {
if (group_type == TOPO_GROUP_TYPE_CHAT && db) {
/* CHAT: добавляем conn только если пир — член именно этого канала */
@ -647,8 +653,9 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
}
}
u_free(chs);
} else {
topo_group_new_conn(group, cqe->conn);
} else if (group_type == TOPO_GROUP_TYPE_UTUN && cqe->conn->links_up) {
struct NODE_CONN_DIRECT* handle = topo_config_handle(g->instance, cqe->peer_node_id);
if (handle) topo_group_join_peer(group, handle);
}
}
entry = entry->hash_next;
@ -762,7 +769,7 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
}
/* Обрабатывает DOWN пира: чистит пути, каскадно удаляет узлы, шлёт withdraw. */
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason) {
if (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; }
if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid instance"); return; }
@ -788,7 +795,9 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
while (pe) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)pe; if (p->conn == conn && p->hop_count >= 2) { uint64_t* hop = (uint64_t*)((uint8_t*)p + sizeof(struct TOPO_NODEPATH)); if (hop[p->hop_count - 1] == conn->peer_node_id) next_hop = hop[p->hop_count - 2]; break; } pe = pe->next; }
if (topo_group_remove_path(nq, conn) == 1) {
uint64_t key = nq->node_id;
if (key != conn->peer_node_id) { topo_recovery_add_node(group, nq, next_hop); cascaded++; }
if (reason == TOPO_REMOVE_TRANSPORT_DOWN && key != conn->peer_node_id) {
topo_recovery_add_node(group, nq, next_hop); cascaded++;
}
nq->conn_presence = 0; nq->conn_up = 0;
topo_fire_nodeinfo_cbk(conn->instance, group, nq);
if (group->group_type != TOPO_GROUP_TYPE_CHAT && rt) route_delete(rt, nq);
@ -842,7 +851,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
topo_group_connect_on_down(group, conn);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d grp=%016llx", conn->log_name, nodes_removed, (unsigned long long)group->group_id);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP peer removed: %s nodes=%d grp=%016llx reason=%d recovery=%d",
conn->log_name, nodes_removed, (unsigned long long)group->group_id, reason, cascaded);
}
@ -1035,7 +1045,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
u_free(new_hop_list); return -1;
}
new_ni = stored;
node_conn_direct_update_node(group->instance, node_id);
if (!is_new_node) topo_nodeq_free_group_fields(group->instance->topo_groups, nodeinfo1);
nodeinfo1->node_id = node_id;
nodeinfo1->subnets = new_subnets;
@ -1119,6 +1128,9 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
if (v4a > 0 || v6a > 0) route_connectivity_probe_node(group->instance, group, nodeinfo1);
}
/* Создание линков может вызвать транспортные callbacks: структура группы уже обновлена. */
node_conn_direct_update_node(group->instance, node_id);
return 0;
}
@ -1154,7 +1166,7 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
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;
size_t max_sz = sizeof(struct TOPOMSG_NODEINFO_PKT) + 4096;
size_t max_sz = sizeof(struct TOPOMSG_NODEINFO_PKT) + TOPO_NODE_WIRE_MAX_SIZE;
uint8_t* p = u_malloc(max_sz);
if (!p) return;
@ -1219,6 +1231,89 @@ static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t
return false;
}
int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle) {
struct ETCP_CONN* conn = node_conn_direct_get_conn(handle);
if (!group || !conn || !conn->links_up) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "join peer requires a connected NCD handle"); return -1;
}
return topo_group_new_conn(group, conn);
}
struct topo_leave {
struct UTUN_INSTANCE* instance;
uint64_t group_id;
struct NODE_CONN_DIRECT* handle;
struct ll_queue* queue;
struct queue_waiter_handle waiter;
};
static void topo_leave_finish(struct topo_leave* leave) {
queue_waiter_cancel(leave->queue, &leave->waiter);
node_conn_direct_close(leave->handle);
u_free(leave);
}
static void topo_leave_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) {
(void)handle;
if (event == NCD_EVENT_UP) return;
struct topo_leave* leave = arg;
DEBUG_WARN(DEBUG_CATEGORY_BGP, "LEAVE_GROUP cancelled: transport event=%d group=%016llx",
event, (unsigned long long)leave->group_id);
/* NCD CLOSED приходит до уничтожения send_input_q; waiter ещё можно отменить. */
topo_leave_finish(leave);
}
static void topo_leave_ready(struct ll_queue* queue, void* arg) {
(void)queue;
struct topo_leave* leave = arg;
struct ETCP_CONN* conn = node_conn_direct_get_conn(leave->handle);
struct TOPO_GROUP* group = topo_groups_find(leave->instance->topo_groups, leave->group_id);
if (group && topo_group_has_sender(group, conn)) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "LEAVE_GROUP superseded by a new membership group=%016llx",
(unsigned long long)leave->group_id);
topo_leave_finish(leave); return;
}
struct ll_entry* packet = ll_alloc_lldgram(sizeof(struct TOPOMSG_JOIN_GROUP));
int result = -1;
if (packet) {
struct TOPOMSG_JOIN_GROUP msg = { .cmd = ETCP_ID_TOPO_ENTRY, .subcmd = TOPO_SUBCMD_LEAVE_GROUP,
.group_id = leave->group_id };
memcpy(packet->dgram, &msg, sizeof(msg)); packet->len = sizeof(msg);
result = etcp_send(conn, packet);
if (result != 0) { queue_dgram_free(packet); queue_entry_free(packet); }
}
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);
topo_leave_finish(leave);
}
int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle) {
struct ETCP_CONN* conn = node_conn_direct_get_conn(handle);
if (!group || !conn) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "leave peer without group/connection"); return -1; }
if (!topo_group_has_sender(group, conn)) return 0;
if (!conn->send_input_q || conn->close_requested) {
topo_group_remove_conn(group, conn, TOPO_REMOVE_LOCAL_LEAVE); return 0;
}
struct topo_leave* leave = u_calloc(1, sizeof(*leave));
if (!leave) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LEAVE_GROUP allocation failed");
topo_group_remove_conn(group, conn, TOPO_REMOVE_LOCAL_LEAVE); return -1;
}
leave->instance = group->instance; leave->group_id = group->group_id; leave->queue = conn->send_input_q;
if (node_conn_direct_open(group->instance, conn->peer_node_id, topo_leave_event, leave, &leave->handle, NULL) == NCD_ERR) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot retain NCD handle for LEAVE_GROUP"); u_free(leave);
topo_group_remove_conn(group, conn, TOPO_REMOVE_LOCAL_LEAVE); return -1;
}
topo_group_remove_conn(group, conn, TOPO_REMOVE_LOCAL_LEAVE);
/* Отправка единственного управляющего пакета после освобождения общей очереди.
Контекст не зависит от срока жизни группы; новый JOIN отменяет устаревший LEAVE. */
if (queue_waiter_wait(leave->queue, &leave->waiter, topo_leave_ready, leave) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot wait to send LEAVE_GROUP"); topo_leave_finish(leave); return -1;
}
return 0;
}
/* Шлёт пиру полную таблицу узлов группы (все узлы, кроме достижимых через него). */
static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (!group || !conn) return;
@ -1262,8 +1357,7 @@ static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_C
static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn) {
if (!instance || !conn || !conn->instance || !instance->topo_groups) return;
/* общий сигнал от пассивной стороны: рестартуем обмен, если мы для этого пира VPN-клиент */
int is_client = conn->instance->config && config_peer_in_clients(conn->instance->config, conn->crypto_ctx.peer_public_key);
if (!is_client && conn->bgp_ready_cbk == NULL) return;
if (!topo_config_handle(instance, conn->peer_node_id)) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_resync: from %s — re-announce groups", conn->log_name);
struct ll_entry* fe = instance->topo_groups->group_list->head;
while (fe) {

12
src/routing_layer/topo_group.h

@ -108,6 +108,7 @@ void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* g
#define TOPO_SUBCMD_ERR_GROUP_MISMATCH 0x0C // ошибка несоответствия типа группы
#define TOPO_SUBCMD_JOIN_GROUP 0x0D // запрос членства в группе
#define TOPO_SUBCMD_RESYNC 0x0E // общий запрос: «я переподключился, заново анонсируй свои группы»
#define TOPO_SUBCMD_LEAVE_GROUP 0x0F // завершить участие, сохранив общий транспорт
#define MAX_HOPS 16
#define BGP_NODES_HASH_SIZE 256
@ -298,7 +299,14 @@ int topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
*
* Вызывается при ETCP on_down.
*/
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
enum topo_group_remove_reason {
TOPO_REMOVE_TRANSPORT_DOWN,
TOPO_REMOVE_LOCAL_LEAVE,
TOPO_REMOVE_REMOTE_LEAVE,
TOPO_REMOVE_MEMBER_INVALID
};
/* Только потеря транспорта запускает recovery каскадно потерянных маршрутов. */
void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason);
/**
* @brief Обрабатывает пакет NODEINFO.
@ -309,6 +317,8 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn);
* @return 0 при успехе
*/
int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from, const uint8_t* data, size_t len);
int topo_group_join_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle);
int topo_group_leave_peer(struct TOPO_GROUP* group, struct NODE_CONN_DIRECT* handle);
/**
* @brief Обрабатывает WITHDRAW.

4
src/routing_layer/topo_node.c

@ -455,7 +455,9 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq,
topo_list_count((struct _topo_head*)ni->v4_addrs) > UINT8_MAX ||
topo_list_count((struct _topo_head*)ni->v6_sock_meta) > UINT8_MAX ||
topo_list_count((struct _topo_head*)ni->v6_addrs) > UINT8_MAX ||
topo_list_count((struct _topo_head*)ni->reality_socks) > UINT8_MAX) {
topo_list_count((struct _topo_head*)ni->reality_socks) > UINT8_MAX ||
(nq->subnets && (topo_list_count((struct _topo_head*)nq->subnets->v4_subnets) > UINT8_MAX ||
topo_list_count((struct _topo_head*)nq->subnets->v6_subnets) > UINT8_MAX))) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "node exceeds wire format limits"); return -1;
}
msg.local_v4_sockets = topo_list_count((struct _topo_head*)ni->v4_sock_meta);

7
src/routing_layer/topo_node.h

@ -117,8 +117,8 @@ struct TOPOMSG_SUBNET4 { uint8_t addr[4]; uint8_t prefix_length; } __attribute
struct TOPOMSG_SUBNET6 { uint8_t addr[16]; uint8_t prefix_length; } __attribute__((packed));
/* Reality-камуфляж сокета: что нужно клиенту для подключения к reality TCP-сокету.
* short_id(8) + server_pubkey(32) + version(3) + server_name(фиксированный 64). */
#define TOPO_REALITY_SNAME_MAX 64
* short_id(8) + server_pubkey(32) + version(3) + server_name(фиксированный 256). */
#define TOPO_REALITY_SNAME_MAX 256
struct TOPOMSG_REALITY_SOCK {
uint8_t socket_id;
uint8_t short_id[8];
@ -160,7 +160,7 @@ struct TOPO_REALITY_SOCK {
};
/* Reality-параметры адреса для персистенции в node_addresses.options (BLOB, фиксированный размер). */
#define TOPO_ADDR_REALITY_OPTS_SIZE 107
#define TOPO_ADDR_REALITY_OPTS_SIZE (8 + 32 + 3 + TOPO_REALITY_SNAME_MAX)
struct TOPO_ADDR_REALITY_OPTS {
uint8_t short_id[8];
uint8_t server_pubkey[32];
@ -253,6 +253,7 @@ uint16_t topo_get_chain_rtt(struct TOPO_GROUP_NODE* nq);
/** Build canonical message for Ed25519 signature: x25519_pubkey || name || client_type || client_activity || addresses */
#define TOPO_SIG_MSG_MAX_SIZE 2048
#define TOPO_NODE_WIRE_MAX_SIZE 16384
int topo_node_build_sig_msg(struct TOPO_NODE* ni, uint8_t* buf, size_t buf_size);
/** Sign self NODEINFO with Ed25519: build_sig_msg + sc_ed25519_sign → ni->x25519_self_sig */

2
src/routing_layer/topo_node_sqlite.c

@ -921,7 +921,7 @@ uint64_t topo_node_sqlite_snapshot_timestamp(sqlite3* db, uint64_t node_id) {
int topo_node_sqlite_snapshot_put(sqlite3* db, struct TOPO_NODE* ni) {
if (!db || topo_node_verify(ni) < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid node snapshot"); return -1; }
uint8_t record[4096];
uint8_t record[TOPO_NODE_WIRE_MAX_SIZE];
struct TOPO_GROUP_NODE nq = {0};
int len = topo_node_serialize(ni, &nq, 0, 0, record, sizeof(record), 0);
if (len < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "snapshot serialization failed"); return -1; }

236
src/transport_layer/etcp_connections.c

@ -2770,53 +2770,111 @@ int init_sockets(struct UTUN_INSTANCE* instance) {
return 0; // All OK
}
/* Создать TCP-линк(и) с REALITY-камуфляжем для [client] с reality=1.
* conn создаётся через NCD (только public_key, без адресов — линки добавляем сами).
* local_srv у link опционален: если задан — это [server] transport=tcp для bind
* на нужный интерфейс (ip); если нет — bind не делается (ОС выбирает source).
* Возвращает NCD-handle (сохранить в config_conn_handles) или NULL. */
struct NODE_CONN_DIRECT* etcp_config_client_reality(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client,
const uint8_t pubkey_bin[SC_PUBKEY_SIZE], uint64_t node_id) {
if (!instance || !client || !pubkey_bin) return NULL;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = client->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
struct NODE_CONN_DIRECT* handle = NULL;
int r = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, NULL);
if (r == NCD_ERR || !handle) {
DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: ncd open failed node=0x%016llx", client->name, (unsigned long long)node_id);
return NULL;
/* Конфигурация владеет только NCD handles; транспортные линки создаёт NCD. */
static void config_client_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) {
struct UTUN_INSTANCE* instance = arg;
if (event == NCD_EVENT_UP) {
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
if (group && topo_group_join_peer(group, handle) < 0)
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "configured peer join failed node=%016llx",
(unsigned long long)node_conn_direct_node_id(handle));
} else if (event == NCD_EVENT_TIMEOUT || event == NCD_EVENT_CLOSED) {
DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "configured peer node=%016llx event=%d",
(unsigned long long)node_conn_direct_node_id(handle), event);
}
struct ETCP_CONN* conn = node_conn_direct_get_conn(handle);
if (!conn) { node_conn_direct_close(handle); return NULL; }
}
int link_count = 0;
for (struct CFG_CLIENT_LINK* cl = client->links; cl; cl = cl->next) {
if (cl->remote_addr.ss_family != AF_INET && cl->remote_addr.ss_family != AF_INET6) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: link without remote addr, skipping", client->name);
continue;
}
struct ETCP_SOCKET* bind_sock = NULL;
if (cl->local_srv) {
static struct NODE_CONN_DIRECT* config_client_open(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client) {
struct TOPO_NODE ni = { .node_name = client->name };
if (sc_hex_to_binary(client->peer_public_key_hex, ni.public_key, SC_PUBKEY_SIZE) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "invalid public key for client %s", client->name); return NULL;
}
ni.node_id = sc_derive_node_id_from_pubkey(ni.public_key);
struct NODE_CONN_DIRECT* handle = NULL;
for (struct CFG_CLIENT_LINK* endpoint = client->links; endpoint; endpoint = endpoint->next) {
struct ETCP_SOCKET* sock = NULL;
if (endpoint->local_srv) {
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (strcmp(cl->local_srv->name, s->name) == 0) { bind_sock = s; break; }
if (!bind_sock) DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: bind socket '%s' not found, using auto", client->name, cl->local_srv->name);
else if (!bind_sock->is_tcp) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: bind socket '%s' is not TCP, ignoring bind", client->name, cl->local_srv->name); bind_sock = NULL; }
if (!strcmp(s->name, endpoint->local_srv->name)) { sock = s; break; }
if (!sock || (client->reality_enabled && !sock->is_tcp)) {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client %s: missing/incompatible bind socket %s",
client->name, endpoint->local_srv->name); goto fail;
}
}
uint8_t protocol = client->reality_enabled ? TOPO_PROTO_TCP | TOPO_PROTO_REALITY :
(sock && sock->is_tcp ? TOPO_PROTO_TCP : TOPO_PROTO_UDP);
struct TOPO_ADDR4 v4 = { .protocol = protocol };
struct TOPO_ADDR6 v6 = { .protocol = protocol };
ni.v4_addrs = NULL; ni.v6_addrs = NULL;
if (endpoint->remote_addr.ss_family == AF_INET) {
const struct sockaddr_in* addr = (const struct sockaddr_in*)&endpoint->remote_addr;
memcpy(v4.addr, &addr->sin_addr, 4); v4.port = ntohs(addr->sin_port); ni.v4_addrs = &v4;
} else if (endpoint->remote_addr.ss_family == AF_INET6) {
const struct sockaddr_in6* addr = (const struct sockaddr_in6*)&endpoint->remote_addr;
memcpy(v6.addr, &addr->sin6_addr, 16); v6.port = ntohs(addr->sin6_port); ni.v6_addrs = &v6;
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client %s: invalid address family", client->name); goto fail;
}
struct ETCP_LINK* tlink = etcp_link_new(conn, bind_sock, &cl->remote_addr, 0);
if (!tlink) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "client %s reality: etcp_link_new failed", client->name); continue; }
tlink->is_tcp = 1;
tlink->reality = client->reality; tlink->reality_set = 1;
etcp_tcp_link_start_connect(tlink, &cl->remote_addr, 0);
link_count++;
DEBUG_INFO(DEBUG_CATEGORY_REALITY, "client %s reality: TCP link %d → %s sn=%s bind=%s",
client->name, tlink->local_link_id, sockaddr_storage_to_str(&cl->remote_addr).str,
client->reality.server_name, bind_sock ? bind_sock->name : "auto");
}
if (link_count == 0) DEBUG_WARN(DEBUG_CATEGORY_REALITY, "client %s reality: no links created", client->name);
struct TOPO_REALITY_SOCK reality = {0};
if (client->reality_enabled) {
memcpy(reality.short_id, client->reality.short_id, sizeof(reality.short_id));
memcpy(reality.server_pubkey, client->reality.server_static_pubkey, sizeof(reality.server_pubkey));
memcpy(reality.version, client->reality.version, sizeof(reality.version));
memcpy(reality.server_name, client->reality.server_name, sizeof(reality.server_name));
ni.reality_socks = &reality;
}
struct NODE_CONN_DIRECT* added = NULL;
if (node_conn_direct_open_node(instance, ni.node_id, NULL, NULL, &added, &ni, sock) == NCD_ERR) goto fail;
if (!handle) handle = added;
else node_conn_direct_close(added); /* один владелец конфигурации, сколько бы адресов ни было */
}
if (!handle && node_conn_direct_open(instance, ni.node_id, NULL, NULL, &handle, NULL) == NCD_ERR) goto fail;
return handle;
fail:
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "cannot configure NCD client %s", client->name);
if (handle) node_conn_direct_close(handle);
return NULL;
}
static void config_handles_close(struct CONFIG_CONN_HANDLE* head) {
while (head) {
struct CONFIG_CONN_HANDLE* next = head->next;
node_conn_direct_close(head->handle); u_free(head); head = next;
}
}
int init_client_connections(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* clients) {
if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "configure clients without instance"); return -1; }
struct CONFIG_CONN_HANDLE* updated = NULL;
int count = 0;
for (struct CFG_CLIENT* client = clients; client; client = client->next) {
struct CONFIG_CONN_HANDLE* owner = u_calloc(1, sizeof(*owner));
if (!owner) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client handle allocation failed"); goto fail; }
owner->handle = config_client_open(instance, client);
if (!owner->handle) { u_free(owner); goto fail; }
owner->node_id = node_conn_direct_node_id(owner->handle);
snprintf(owner->name, sizeof(owner->name), "%s", client->name);
owner->next = updated; updated = owner; count++;
}
struct CONFIG_CONN_HANDLE* old = instance->config_conn_handles;
instance->config_conn_handles = updated;
struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups);
for (struct CONFIG_CONN_HANDLE* owner = old; owner; owner = owner->next) {
int retained = 0;
for (struct CONFIG_CONN_HANDLE* next = updated; next; next = next->next)
if (next->node_id == owner->node_id) { retained = 1; break; }
if (!retained && group && node_conn_direct_get_conn(owner->handle)) topo_group_leave_peer(group, owner->handle);
}
for (struct CONFIG_CONN_HANDLE* owner = updated; owner; owner = owner->next) {
node_conn_direct_set_callback(owner->handle, config_client_event, instance);
}
config_handles_close(old);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "configured %d NCD client handles", count);
return 0;
fail:
config_handles_close(updated);
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "NCD client configuration rejected; previous handles retained");
return -1;
}
// Инициализация сети: listen-сокеты (если ещё нет) + client-соединения через NCD из конфига.
@ -2843,99 +2901,7 @@ int init_connections(struct UTUN_INSTANCE* instance) {
etcp_keepalive_register(instance);
// Initialize clients via node_conn_direct
struct CFG_CLIENT* client = config->clients;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "init_connections: clients=%p total_conns=%d", config->clients, queue_entry_count(instance->connections));
while (client) {
if (strlen(client->peer_public_key_hex) == 0) {
DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "no peer public key configured for client %s", client->name);
client = client->next; continue;
}
uint8_t pubkey_bin[SC_PUBKEY_SIZE];
if (sc_hex_to_binary(client->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid peer pubkey hex for client %s", client->name);
client = client->next; continue;
}
uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "client %s node_id=0x%016llx", client->name, (unsigned long long)node_id);
if (client->reality_enabled) {
struct NODE_CONN_DIRECT* rhandle = etcp_config_client_reality(instance, client, pubkey_bin, node_id);
if (rhandle) {
struct CONFIG_CONN_HANDLE* ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, client->name, MAX_CONN_NAME_LEN - 1);
ch->handle = rhandle; ch->next = instance->config_conn_handles;
instance->config_conn_handles = ch;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s saved reality handle to config_conn_handles", client->name); }
}
client = client->next; continue;
}
struct NODE_CONN_DIRECT* handle = NULL;
struct ETCP_CONN* conn = NULL;
for (struct CFG_CLIENT_LINK* cl = client->links; cl; cl = cl->next) {
if (cl->local_srv && cl->local_srv->transport) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "client %s TCP transport not yet supported via NCD, skipping", client->name);
continue;
}
struct ETCP_SOCKET* sock = NULL;
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (cl->local_srv && strcmp(cl->local_srv->name, s->name) == 0) { sock = s; break; }
if (!sock) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s link",
cl->local_srv ? cl->local_srv->name : "?", client->name);
continue;
}
if (!handle) {
struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = client->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
if (cl->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&cl->remote_addr;
v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4);
v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v4_addrs = &v4_addr;
} else if (cl->remote_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&cl->remote_addr;
v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16);
v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v6_addrs = &v6_addr;
}
int r = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock);
if (r == NCD_ERR || !handle) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd open failed for client %s", client->name);
break;
}
conn = node_conn_direct_get_conn(handle);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s ncd handle=%p conn=%p %s sock=%s",
client->name, handle, conn, r == NCD_NEW ? "NEW" : "REUSED", sock->name);
} else {
struct ETCP_LINK* link = etcp_link_new(conn, sock, &cl->remote_addr, 0);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "client %s added link link_id=%d sock=%s", client->name, link ? link->local_link_id : -1, sock->name);
}
}
if (handle) {
struct CONFIG_CONN_HANDLE* ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, client->name, MAX_CONN_NAME_LEN - 1);
ch->handle = handle; ch->next = instance->config_conn_handles;
instance->config_conn_handles = ch;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "client %s saved handle to config_conn_handles", client->name); }
}
client = client->next;
}
// If there are clients configured but no connections created, that's an error
// If there are no clients (server-only mode), 0 connections is OK (server will accept incoming)
int total_conns = queue_entry_count(instance->connections);
if (total_conns == 0 && config->clients != NULL) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Clients configured but no connections initialized");
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized %d connections", total_conns);
if (init_client_connections(instance, config->clients) < 0) return -1;
// Return 1 if there was a partial socket initialization error
if (socket_result == 1) return 1;
return 0;

6
src/transport_layer/etcp_connections.h

@ -354,10 +354,8 @@ int init_sockets(struct UTUN_INSTANCE* instance);
// Создаёт listen-сокеты и client connections из конфига
int init_connections(struct UTUN_INSTANCE* instance);
// Создать TCP-линк(и) с REALITY-камуфляжем для [client] с reality=1.
// local_srv у link опционален (bind на интерфейс); возвращает NCD-handle или NULL.
struct NODE_CONN_DIRECT* etcp_config_client_reality(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client,
const uint8_t pubkey_bin[SC_PUBKEY_SIZE], uint64_t node_id);
// Заменяет принадлежащие конфигурации NCD handles. Новые приобретаются до освобождения старых.
int init_client_connections(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* clients);
// SOCKET FUNCTIONS
// добавляет новый версер (сокет для приёма и отправки кодограмм. обслуживает много подключений)

44
src/transport_layer/node_conn_direct.c

@ -195,8 +195,19 @@ static int ncd_same_address(const struct sockaddr_storage* a, const struct socka
static int ncd_add_tcp_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* sock,
struct sockaddr_storage* addr, const struct TOPO_REALITY_SOCK* reality) {
for (struct ETCP_LINK* link = conn->links; link; link = link->next)
if (link->is_tcp && !link->is_server && link->conn == sock && ncd_same_address(&link->remote_addr, addr)) return 0;
for (struct ETCP_LINK* link = conn->links; link; link = link->next) {
if (!link->is_tcp || link->is_server || link->conn != sock || !ncd_same_address(&link->remote_addr, addr)) continue;
struct reality_client_config before; memcpy(&before, &link->reality, sizeof(before));
int had_reality = link->reality_set;
if (reality) topo_node_apply_reality(link, reality);
else { memset(&link->reality, 0, sizeof(link->reality)); link->reality_set = 0; }
if (had_reality != link->reality_set || memcmp(&before, &link->reality, sizeof(before))) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] TCP parameters updated node=%016llx link=%d up=%d",
(unsigned long long)conn->peer_node_id, link->local_link_id, link->link_status);
if (!link->link_status) etcp_tcp_link_start_reconnect(link);
}
return 0;
}
struct ETCP_LINK* link = etcp_link_new(conn, sock, addr, 0);
if (!link) { DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] TCP link allocation failed"); return -1; }
topo_node_apply_reality(link, reality);
@ -243,12 +254,23 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
uint8_t tnat = m ? m->nat_type : NAT_TYPE_UNKNOWN;
if (a->protocol & TOPO_PROTO_TCP) {
if ((a->protocol & TOPO_PROTO_REALITY) && !topo_node_find_reality_sock(ni, a->socket_id)) {
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] missing REALITY metadata node=%016llx socket=%u",
(unsigned long long)nid, a->socket_id); continue;
}
int tcp_sockets = 0;
for (int vi = 0; vi < nviews; vi++) {
const struct sock_view* v = &views[vi];
if (!v->is_tcp || (specific_sock && v->sock != specific_sock) || !sock_match(v, a->type, tcfg, tnat, tip)) continue;
if (!v->is_tcp) continue;
tcp_sockets++;
if (specific_sock ? v->sock != specific_sock : !sock_match(v, a->type, tcfg, tnat, tip)) continue;
int added = ncd_add_tcp_link(conn, v->sock, &sa, topo_node_find_reality_sock(ni, a->socket_id));
if (added > 0) link_count += added;
}
if (!specific_sock && !tcp_sockets) {
int added = ncd_add_tcp_link(conn, NULL, &sa, topo_node_find_reality_sock(ni, a->socket_id));
if (added > 0) link_count += added;
}
}
if (a->protocol & TOPO_PROTO_UDP) {
if (specific_sock && !specific_sock->is_tcp && specific_sock->local_addr.ss_family == AF_INET) {
@ -274,16 +296,26 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
memcpy(&sin6.sin6_addr, a->addr, 16); sin6.sin6_port = htons(a->port);
if (a->protocol & TOPO_PROTO_TCP) {
if ((a->protocol & TOPO_PROTO_REALITY) && !topo_node_find_reality_sock(ni, a->socket_id)) {
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] missing REALITY metadata node=%016llx socket=%u",
(unsigned long long)nid, a->socket_id); continue;
}
int tcp_sockets = 0;
if (is_ll) { struct ETCP_SOCKET* sv = inst->etcp_sockets;
while (sv) { if (sv->local_addr.ss_family == AF_INET6 && sv->netif_index) { sin6.sin6_scope_id = sv->netif_index; break; } sv = sv->next; } }
struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6));
for (struct ETCP_SOCKET* s = inst->etcp_sockets; s; s = s->next) {
if (!s->is_tcp || (specific_sock && s != specific_sock) || s->local_addr.ss_family != AF_INET6 ||
s->type == CFG_SERVER_TYPE_PRIVATE) continue;
if (!s->is_tcp || s->local_addr.ss_family != AF_INET6) continue;
tcp_sockets++;
if (specific_sock ? s != specific_sock : s->type == CFG_SERVER_TYPE_PRIVATE) continue;
if (is_ll) { sin6.sin6_scope_id = s->netif_index; memcpy(&sa, &sin6, sizeof(sin6)); }
int added = ncd_add_tcp_link(conn, s, &sa, topo_node_find_reality_sock(ni, a->socket_id));
if (added > 0) link_count += added;
}
if (!specific_sock && !tcp_sockets && !is_ll) {
int added = ncd_add_tcp_link(conn, NULL, &sa, topo_node_find_reality_sock(ni, a->socket_id));
if (added > 0) link_count += added;
}
}
if (a->protocol & TOPO_PROTO_UDP) {
if (specific_sock && !specific_sock->is_tcp && specific_sock->local_addr.ss_family == AF_INET6) {
@ -971,7 +1003,7 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
}
struct TOPO_NODE* known = ncd_lookup_node(inst, node_id);
if (ni->timestamp && (!known || ni->timestamp > known->timestamp)) {
uint8_t wire[4096];
uint8_t wire[TOPO_NODE_WIRE_MAX_SIZE];
struct TOPO_GROUP_NODE nq = {0};
struct TOPO_GROUP group = { .instance = inst };
struct TOPO_NODE* copy = NULL;

64
src/utun_instance.c

@ -1009,64 +1009,9 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc
}
os = next;
}
// Clients [client:] — close old config handles, rebuild from new config
struct CONFIG_CONN_HANDLE* ch = instance->config_conn_handles;
while (ch) { struct CONFIG_CONN_HANDLE* next = ch->next; node_conn_direct_close(ch->handle); u_free(ch); ch = next; }
instance->config_conn_handles = NULL;
for (struct CFG_CLIENT *nc = new_config->clients; nc; nc = nc->next) {
if (strlen(nc->peer_public_key_hex) == 0) { DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "client %s no peer pubkey", nc->name); continue; }
uint8_t pubkey_bin[SC_PUBKEY_SIZE];
if (sc_hex_to_binary(nc->peer_public_key_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "invalid pubkey hex client %s", nc->name); continue; }
uint64_t node_id = sc_derive_node_id_from_pubkey(pubkey_bin);
if (nc->reality_enabled) {
struct NODE_CONN_DIRECT* rhandle = etcp_config_client_reality(instance, nc, pubkey_bin, node_id);
if (rhandle) {
ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, nc->name, MAX_CONN_NAME_LEN - 1);
ch->handle = rhandle; ch->next = instance->config_conn_handles; instance->config_conn_handles = ch; }
}
continue;
}
struct NODE_CONN_DIRECT* handle = NULL;
struct ETCP_CONN* conn = NULL;
for (struct CFG_CLIENT_LINK *nl = nc->links; nl; nl = nl->next) {
struct ETCP_SOCKET* sock = NULL;
for (struct ETCP_SOCKET* s = instance->etcp_sockets; s; s = s->next)
if (nl->local_srv && strcmp(nl->local_srv->name, s->name) == 0) { sock = s; break; }
if (!sock) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "no socket for server '%s' client %s", nl->local_srv ? nl->local_srv->name : "?", nc->name); continue; }
if (!handle) {
struct TOPO_ADDR4 v4_addr; struct TOPO_ADDR6 v6_addr;
struct TOPO_NODE ni_tmp; memset(&ni_tmp, 0, sizeof(ni_tmp));
ni_tmp.node_id = node_id; ni_tmp.node_name = nc->name;
memcpy(ni_tmp.public_key, pubkey_bin, SC_PUBKEY_SIZE);
if (nl->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&nl->remote_addr;
v4_addr = (struct TOPO_ADDR4){0}; memcpy(v4_addr.addr, &sin->sin_addr.s_addr, 4);
v4_addr.port = ntohs(sin->sin_port); v4_addr.type = TOPO_ADDR_INTERFACE; v4_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v4_addrs = &v4_addr;
} else if (nl->remote_addr.ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)&nl->remote_addr;
v6_addr = (struct TOPO_ADDR6){0}; memcpy(v6_addr.addr, &sin6->sin6_addr, 16);
v6_addr.port = ntohs(sin6->sin6_port); v6_addr.type = TOPO_ADDR_INTERFACE; v6_addr.protocol = TOPO_PROTO_UDP;
ni_tmp.v6_addrs = &v6_addr;
}
int rc = node_conn_direct_open_node(instance, node_id, NULL, NULL, &handle, &ni_tmp, sock);
if (rc == NCD_ERR || !handle) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "ncd reload open failed client %s", nc->name); break; }
conn = node_conn_direct_get_conn(handle);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "reload: client %s ncd %s node=0x%016llx", nc->name, rc == NCD_NEW ? "NEW" : "REUSED", (unsigned long long)node_id);
} else {
etcp_link_new(conn, sock, &nl->remote_addr, 0);
}
}
if (handle) {
ch = u_calloc(1, sizeof(*ch));
if (ch) { ch->node_id = node_id; strncpy(ch->name, nc->name, MAX_CONN_NAME_LEN - 1);
ch->handle = handle; ch->next = instance->config_conn_handles; instance->config_conn_handles = ch; }
}
}
int client_result = init_client_connections(instance, new_config->clients);
if (client_result < 0)
DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "reload: failed to replace client handles; previous connections retained");
// Reload networks: clear and repopulate
if (instance->networks) {
struct ll_entry *entry;
@ -1103,7 +1048,8 @@ struct UTUN_INSTANCE *utun_instance_reload(struct UTUN_INSTANCE *instance, struc
// routing_create(instance); // commented per request
if (instance->config) free_config(instance->config);
instance->config = new_config;
DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Selective partial reload completed (only changed updated, unchanged untouched, [debug] applied)");
if (client_result < 0) DEBUG_WARN(DEBUG_CATEGORY_CONFIG, "Reload completed with client errors; previous client handles retained");
else DEBUG_INFO(DEBUG_CATEGORY_CONFIG, "Selective partial reload completed (only changed updated, unchanged untouched, [debug] applied)");
return instance;
}

9
tests/Makefile.am

@ -66,6 +66,7 @@ check_PROGRAMS = \
test_node_conn_direct \
test_group_ownership \
test_node_snapshot \
test_ncd_config \
test_db_sync \
test_merkle_sync \
test_merkle_protocol \
@ -345,11 +346,11 @@ test_route_lib_SOURCES = test_route_lib.c
test_route_lib_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_route_lib_LDADD = $(LIBUTUN) $(COMMON_LIBS)
test_bgp_route_exchange_SOURCES = test_bgp_route_exchange.c
test_bgp_route_exchange_SOURCES = test_bgp_route_exchange.c test_connection_loss.h
test_bgp_route_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_bgp_route_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_bgp_triangle_SOURCES = test_bgp_triangle.c
test_bgp_triangle_SOURCES = test_bgp_triangle.c test_connection_loss.h
test_bgp_triangle_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_bgp_triangle_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
@ -397,6 +398,10 @@ 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)
test_ncd_config_SOURCES = test_ncd_config.c
test_ncd_config_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_ncd_config_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_db_sync_SOURCES = test_db_sync.c
test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

31
tests/test_bgp_route_exchange.c

@ -42,7 +42,7 @@
#include "routing.h"
#include "route_lib.h"
#include "topo_group.h"
#include "dummynet.h"
#include "test_connection_loss.h"
#include "../src/tun_if.h"
#include "secure_channel.h"
#include "../lib/u_async.h"
@ -65,9 +65,7 @@ static int test_phase = 0; // 0=running, 1=success, 2=failure
static void* timeout_id = NULL;
// dummynet filters per link
static struct dummynet_filter* df_ab1 = NULL;
static struct dummynet_filter* df_ab2 = NULL;
static struct dummynet_filter* df_bc = NULL;
static struct test_connection_loss loss_ab, loss_bc;
static char temp_dir[] = "/tmp/utun_bgp_mesh_XXXXXX";
static char config_a[256], config_b[256], config_c[256];
@ -337,10 +335,8 @@ int main(void) {
ASSERT(ab_link1 && ab_link2, "A↔B links not found");
ASSERT(bc_link, "B↔C link not found");
// Create dummynet filters
df_ab1 = dummynet_filter_create(ua); ASSERT(df_ab1, "df_ab1");
df_ab2 = dummynet_filter_create(ua); ASSERT(df_ab2, "df_ab2");
df_bc = dummynet_filter_create(ua); ASSERT(df_bc, "df_bc");
test_loss_init(&loss_ab, inst_a, g_node_id_b);
test_loss_init(&loss_bc, inst_b, g_node_id_c);
timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout_cb, "test_timeout");
@ -352,8 +348,7 @@ int main(void) {
// --- Phase 2: Kill one A↔B link, verify routes stay ---
PHASE("2: Kill one A↔B link");
dummynet_filter_attach(df_ab1, ab_link1);
dummynet_filter_set_loss(df_ab1, 1000);
test_loss_set(&loss_ab, 1, ab_link1, NULL);
// Wait for link-down detection by keepalive, then verify routes survive (link2 still up)
for (int i = 0; i < 200 && test_phase == 0 && ab_link1->link_status; i++) uasync_poll(ua, POLL_INTERVAL_MS);
// Small margin after link-down before checks
@ -363,33 +358,30 @@ int main(void) {
// --- Phase 3: Restore the link ---
PHASE("3: Restore one A↔B link");
dummynet_filter_set_loss(df_ab1, 0);
test_loss_set(&loss_ab, 0, NULL, NULL);
ASSERT(wait_for("BGP recovery after ab1 restored", cond_bgp_a_to_c, PHASE_TIMEOUT_TB), "bgp A→C recovery");
// --- Phase 4: Kill B↔C link ---
PHASE("4: Kill B↔C link");
dummynet_filter_attach(df_bc, bc_link);
dummynet_filter_set_loss(df_bc, 1000);
test_loss_set(&loss_bc, 1, NULL, NULL);
ASSERT(wait_for("A route to C removed", cond_bgp_a_to_c_gone, PHASE_TIMEOUT_TB), "bgp A→C gone");
ASSERT(wait_for("C route to A removed", cond_bgp_c_to_a_gone, PHASE_TIMEOUT_TB), "bgp C→A gone");
// --- Phase 5: Restore B↔C ---
PHASE("5: Restore B↔C link");
dummynet_filter_set_loss(df_bc, 0);
test_loss_set(&loss_bc, 0, NULL, NULL);
ASSERT(wait_for("BGP A→C restored", cond_bgp_a_to_c, PHASE_TIMEOUT_TB), "bgp A→C restored");
ASSERT(wait_for("BGP C→A restored", cond_bgp_c_to_a, PHASE_TIMEOUT_TB), "bgp C→A restored");
// --- Phase 6: Kill both A↔B links ---
PHASE("6: Kill both A↔B links");
dummynet_filter_set_loss(df_ab1, 1000);
dummynet_filter_attach(df_ab2, ab_link2);
dummynet_filter_set_loss(df_ab2, 1000);
test_loss_set(&loss_ab, 1, NULL, NULL);
ASSERT(wait_for("A route to C gone (both links)", cond_bgp_a_to_c_gone, PHASE_TIMEOUT_TB), "bgp A→C gone 2");
ASSERT(wait_for("C route to A gone (both links)", cond_bgp_c_to_a_gone, PHASE_TIMEOUT_TB), "bgp C→A gone 2");
// --- Phase 7: Restore one A↔B link (ab1 only, ab2 stays dead) ---
PHASE("7: Restore one A↔B link");
dummynet_filter_set_loss(df_ab1, 0);
test_loss_set(&loss_ab, 1, NULL, ab_link1);
ASSERT(wait_for("BGP A→C restored via one link", cond_bgp_a_to_c, PHASE_TIMEOUT_TB), "bgp A→C restored 2");
ASSERT(wait_for("BGP C→A restored via one link", cond_bgp_c_to_a, PHASE_TIMEOUT_TB), "bgp C→A restored 2");
@ -397,8 +389,7 @@ int main(void) {
test_phase = 1;
// Cleanup
dummynet_filter_detach(df_ab1); dummynet_filter_detach(df_ab2); dummynet_filter_detach(df_bc);
dummynet_filter_destroy(df_ab1); dummynet_filter_destroy(df_ab2); dummynet_filter_destroy(df_bc);
test_loss_destroy(&loss_ab); test_loss_destroy(&loss_bc);
if (timeout_id) uasync_cancel_timeout(ua, timeout_id);
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); }
if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); }

27
tests/test_bgp_triangle.c

@ -45,7 +45,7 @@
#include "route_lib.h"
#include "topo_group.h"
#include "topo_node.h"
#include "dummynet.h"
#include "test_connection_loss.h"
#include "../src/tun_if.h"
#include "secure_channel.h"
#include "../lib/u_async.h"
@ -71,8 +71,7 @@ static struct UASYNC* ua = NULL;
static int test_phase = 0;
static void* timeout_id = NULL;
static struct dummynet_filter* df_ab = NULL;
static struct dummynet_filter* df_ac = NULL;
static struct test_connection_loss loss_ab, loss_ac;
/* ================================================================
* Config generator
@ -352,8 +351,8 @@ int main(void) {
struct ETCP_LINK* ac_link = find_client_link(inst_a, "to_c");
ASSERT(ab_link && ac_link, "A links not found");
df_ab = dummynet_filter_create(ua); ASSERT(df_ab, "df_ab");
df_ac = dummynet_filter_create(ua); ASSERT(df_ac, "df_ac");
test_loss_init(&loss_ab, inst_a, g_node_id_b);
test_loss_init(&loss_ac, inst_a, g_node_id_c);
timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout_cb, "test_timeout");
@ -397,8 +396,7 @@ int main(void) {
* Phase 2: Kill A–C — C/D/E only through B
* ================================================================ */
PHASE("2: Kill A–C — C/D/E via B only");
dummynet_filter_attach(df_ac, ac_link);
dummynet_filter_set_loss(df_ac, 1000);
test_loss_set(&loss_ac, 1, NULL, NULL);
ASSERT(wait_for("C: 1 path via B after A-C kill", cond_c_one_path_B, PHASE_TIMEOUT_TB), "C paths not 1");
{
int pc = node_path_count(inst_a, g_node_id_c);
@ -418,7 +416,7 @@ int main(void) {
* Fix#1: path [C] added despite old version.
* ================================================================ */
PHASE("3: Restore A–C");
dummynet_filter_set_loss(df_ac, 0);
test_loss_set(&loss_ac, 0, NULL, NULL);
ASSERT(wait_for("C: >=2 paths after A-C restore", cond_c_two_paths, PHASE_TIMEOUT_TB), "C paths not >=2");
{
int cp = node_path_count(inst_a, g_node_id_c);
@ -429,8 +427,7 @@ int main(void) {
* Phase 4: Kill A–B — B/E survive via C relay, C/D via A-C
* ================================================================ */
PHASE("4: Kill A–B — B/E survive via C");
dummynet_filter_attach(df_ab, ab_link);
dummynet_filter_set_loss(df_ab, 1000);
test_loss_set(&loss_ab, 1, NULL, NULL);
/* After A-B kill: direct paths removed, B/C/E survive via C relay */
ASSERT(wait_for("A-B paths removed", cond_c_one_path, PHASE_TIMEOUT_TB), "A-B paths not removed");
{
@ -450,7 +447,7 @@ int main(void) {
* Phase 5: Kill A–C — A fully isolated (A-B already dead)
* ================================================================ */
PHASE("5: Kill A–C — A isolated");
dummynet_filter_set_loss(df_ac, 1000);
test_loss_set(&loss_ac, 1, NULL, NULL);
ASSERT(wait_for("C gone from A", cond_c_gone_a, PHASE_TIMEOUT_TB), "C still in A");
ASSERT(wait_for("D gone from A", cond_d_gone_a, PHASE_TIMEOUT_TB), "D still in A");
ASSERT(!peer_in_nodes(inst_a, g_node_id_b), "B should be gone");
@ -463,8 +460,8 @@ int main(void) {
* Phase 6: Restore A–C + A–B — full recovery
* ================================================================ */
PHASE("6: Full restore — both links back");
dummynet_filter_set_loss(df_ac, 0);
dummynet_filter_set_loss(df_ab, 0);
test_loss_set(&loss_ac, 0, NULL, NULL);
test_loss_set(&loss_ab, 0, NULL, NULL);
ASSERT(wait_for("D in A after restore", cond_d_in_a, 2 * PHASE_TIMEOUT_TB), "D not back");
ASSERT(wait_for("E in A after restore", cond_e_in_a, PHASE_TIMEOUT_TB), "E not back");
ASSERT(wait_for("C: >=2 paths after restore", cond_c_two_paths, PHASE_TIMEOUT_TB), "C paths not >=2");
@ -486,8 +483,8 @@ int main(void) {
test_phase = 1;
cleanup:
if (df_ab) { dummynet_filter_detach(df_ab); dummynet_filter_destroy(df_ab); }
if (df_ac) { dummynet_filter_detach(df_ac); dummynet_filter_destroy(df_ac); }
test_loss_destroy(&loss_ab);
test_loss_destroy(&loss_ac);
if (timeout_id) uasync_cancel_timeout(ua, timeout_id);
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); }
if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); }

65
tests/test_connection_loss.h

@ -0,0 +1,65 @@
#ifndef TEST_CONNECTION_LOSS_H
#define TEST_CONNECTION_LOSS_H
#include "etcp_api.h"
#include "etcp_connections.h"
#include "utun_instance.h"
#include "../lib/debug_config.h"
/* Разрыв пары узлов охватывает и новые линки из свежего NODEINFO.
* only/spared позволяют отдельно проверить потерю одного линка и работу
* единственного оставшегося линка. Контекст живёт до loss_destroy(). */
struct test_connection_loss {
struct UTUN_INSTANCE* instance;
uint64_t peer_id;
struct ETCP_LINK* only;
struct ETCP_LINK* spared;
int blocked;
unsigned links;
unsigned long long dropped;
};
static ssize_t test_loss_send(socket_t fd, const void* buf, size_t len, const struct sockaddr* addr,
socklen_t addr_len, struct ETCP_LINK* link, void* arg) {
struct test_connection_loss* loss = arg;
if (loss->blocked && (!loss->only || loss->only == link) && loss->spared != link) {
loss->dropped++;
return (ssize_t)len;
}
return socket_sendto(fd, buf, len, addr, addr_len);
}
static void test_loss_link(struct ETCP_CONN* conn, struct ETCP_LINK* link, int old_state, int old_status, void* arg) {
(void)old_state; (void)old_status;
struct test_connection_loss* loss = arg;
if (conn->peer_node_id != loss->peer_id || link->send_hook == test_loss_send) return;
link->send_hook = test_loss_send; link->send_hook_ctx = loss; loss->links++;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "test partition: peer=%016llx attached link=%p total=%u blocked=%d",
(unsigned long long)loss->peer_id, (void*)link, loss->links, loss->blocked);
}
static void test_loss_init(struct test_connection_loss* loss, struct UTUN_INSTANCE* instance, uint64_t peer_id) {
loss->instance = instance; loss->peer_id = peer_id;
etcp_add_link_status_cbk(instance, test_loss_link, loss);
struct ETCP_CONN* conn = instance_find_conn(instance, peer_id);
for (struct ETCP_LINK* link = conn ? conn->links : NULL; link; link = link->next)
test_loss_link(conn, link, 0, 0, loss);
}
static void test_loss_set(struct test_connection_loss* loss, int blocked, struct ETCP_LINK* only, struct ETCP_LINK* spared) {
loss->blocked = blocked; loss->only = only; loss->spared = spared;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "test partition: peer=%016llx blocked=%d only=%p spared=%p links=%u dropped=%llu",
(unsigned long long)loss->peer_id, blocked, (void*)only, (void*)spared, loss->links, loss->dropped);
}
static void test_loss_destroy(struct test_connection_loss* loss) {
if (!loss->instance) return;
etcp_remove_link_status_cbk(loss->instance, test_loss_link, loss);
struct ETCP_CONN* conn = instance_find_conn(loss->instance, loss->peer_id);
for (struct ETCP_LINK* link = conn ? conn->links : NULL; link; link = link->next) {
if (link->send_hook_ctx != loss) continue;
link->send_hook = NULL; link->send_hook_ctx = NULL;
}
loss->instance = NULL;
}
#endif

41
tests/test_etcp_connect.c

@ -11,6 +11,7 @@
#include "etcp.h"
#include "etcp_connections.h"
#include "etcp_api.h"
#include "node_conn_direct.h"
#include "../src/config_parser.h"
#include "../src/config_updater.h"
#include "../src/utun_instance.h"
@ -34,6 +35,7 @@ static void* ttimer = NULL;
static char tdir[] = "/tmp/utun_ec_XXXXXX";
static char ca[256], cb[256];
static int pa = 0, pb = 0;
static struct NODE_CONN_DIRECT *chat_a = NULL, *chat_b = NULL;
static int wf(const char* p, const char* f, ...) {
va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1;
@ -60,6 +62,7 @@ static void test3(void* arg); static void test4(void* arg);
static void test5(void* arg); static void test6(void* arg);
static void test7(void* arg); static void test8(void* arg);
static void test9(void* arg); static void test10(void* arg);
static void test11(void* arg);
/* ----- callback state per test ----- */
static volatile int cb_type = 0, cb_conn_ok = 0, cb_count = 0;
@ -197,6 +200,42 @@ static void test10(void* arg) {
if (cb_count != 1) { fail("test9: expected exactly 1 callback"); return; }
if (cb_type != ETCP_CONNECT_LATE) { fail("test9: expected LATE type"); return; }
fprintf(stderr, " OK: only LATE fired\n"); fflush(stderr);
if (node_conn_direct_open(g_a, nid_b, NULL, NULL, &chat_a, NULL) != NCD_REUSED ||
node_conn_direct_open(g_b, nid_a, NULL, NULL, &chat_b, NULL) != NCD_REUSED) {
fail("cannot retain shared connection for chat"); return;
}
struct ETCP_CONN* held = node_conn_direct_get_conn(chat_a);
struct ll_queue* queue = held->send_input_q;
queue->callback_suspended = 1;
struct ll_entry* packet = ll_alloc_lldgram(2);
if (!packet) { fail("cannot allocate backpressure fixture"); return; }
packet->dgram[0] = ETCP_ID_TOPO_ENTRY; packet->dgram[1] = TOPO_SUBCMD_RESYNC; packet->len = 2;
if (etcp_send(held, packet) != 0) { fail("cannot queue backpressure fixture"); return; }
if (init_client_connections(g_a, NULL) < 0) { fail("cannot remove UTUN config owner"); return; }
if (!queue->waiter_head) { fail("LEAVE_GROUP did not wait for backpressure"); return; }
queue_resume_callback(queue);
uasync_call_soon(ua, NULL, test11);
}
static void test11(void* arg) {
(void)arg;
struct ETCP_CONN* a = node_conn_direct_get_conn(chat_a);
struct ETCP_CONN* b = node_conn_direct_get_conn(chat_b);
if (!a || !b || !a->links_up || !b->links_up || a->fin_wait || b->fin_wait) {
fail("removing UTUN membership disrupted another NCD owner"); return;
}
if (topo_node_find_by_id(topo_groups_get_default(g_a->topo_groups), nid_b) ||
topo_node_find_by_id(topo_groups_get_default(g_b->topo_groups), nid_a)) {
uasync_set_timeout(ua, 50, NULL, test11, "leave_group"); return;
}
fprintf(stderr, " OK: UTUN peers detached on both sides, shared transport remains UP\n");
/* Shutdown обязан отменить ещё не отправленный LEAVE и освободить его waiter/handle. */
struct ll_queue* queue = a->send_input_q;
queue->callback_suspended = 1;
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;
}
fprintf(stderr, "=== ALL PASSED ===\n"); fflush(stderr);
result = 2;
}
@ -248,6 +287,8 @@ int main(void) {
fprintf(stderr, "final result=%d\n", result); fflush(stderr);
done:
if (ttimer) uasync_cancel_timeout(ua, ttimer);
if (chat_a) node_conn_direct_close(chat_a);
if (chat_b) node_conn_direct_close(chat_b);
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); }
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); }
if (ua) {

18
tests/test_etcp_router.c

@ -627,9 +627,11 @@ static void timeout_cb(void* arg) {
if (g_mon_id) { uasync_cancel_timeout(ua, g_mon_id); g_mon_id = NULL; }
}
static void tcp_connected(void* arg, struct ETCP_CONN* conn, int type) {
(void)arg; (void)type;
if (!conn) { printf("[FAIL] TCP connect callback\n"); g_done = -1; }
static void tcp_connected(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) {
(void)arg;
if (event == NCD_EVENT_UP && topo_group_join_peer(topo_groups_get_default(cli->topo_groups), handle) < 0) {
printf("[FAIL] TCP group join\n"); g_done = -1;
}
}
// ======================== Main ========================
@ -686,17 +688,9 @@ int main(void) {
ni.node_id = server_node_id; memcpy(ni.public_key,srv->my_keys.public_key,SC_PUBKEY_SIZE);
addr.addr[0]=127; addr.addr[3]=1; addr.port=g_srv_port; addr.protocol=TOPO_PROTO_TCP; addr.type=TOPO_ADDR_INTERFACE;
ni.v4_addrs=&addr;
if (node_conn_direct_open_node(cli,server_node_id,NULL,NULL,&g_tcp_handle,&ni,NULL) == NCD_ERR) {
if (node_conn_direct_open_node(cli,server_node_id,tcp_connected,NULL,&g_tcp_handle,&ni,NULL) == NCD_ERR) {
printf("[FAIL] TCP open\n"); goto done;
}
struct TOPO_NODE* known = u_calloc(1,sizeof(*known));
if (!known) { printf("[FAIL] TCP registry allocation\n"); goto done; }
known->node_id = ni.node_id; memcpy(known->public_key,ni.public_key,SC_PUBKEY_SIZE);
if (!topo_node_registry_store(cli->topo_groups,known)) { printf("[FAIL] TCP registry store\n"); goto done; }
struct TOPO_GROUP_NODE node = {0}; node.node_id = server_node_id;
if (etcp_connect(cli,&node,tcp_connected,NULL,ETCP_CONNECT_BGP_READY) < 0) {
printf("[FAIL] TCP BGP connect\n"); goto done;
}
}
// Create dummynet between client and server

48
tests/test_etcp_router_reconnect.c

@ -66,7 +66,6 @@ static uint32_t g_total_sent_phase1 = 0, g_total_sent_phase3 = 0;
static uint32_t g_phase_sent = 0, g_seq = 0;
static int g_tick_counter = 0;
static uint32_t g_rcvd_total = 0, g_expected_seq = 0;
static int g_bgp_ready_mask = 0;
static int g_loop_cnt = 0;
static void hex_encode(const uint8_t* bin, int len, char* out) {
@ -78,23 +77,6 @@ static void gen_payload(uint32_t seq, uint8_t* buf, int len) {
for (int i=0;i<len;i++) buf[i]=(uint8_t)((i^seq^0xA5)&0xFF);
}
static struct TOPO_GROUP_NODE* mknode(struct UTUN_INSTANCE* inst, uint64_t nid, const uint8_t pubkey[32],
uint8_t a, uint8_t b, uint8_t c, uint8_t d, uint16_t port) {
struct TOPO_GROUP_NODE* nq = u_calloc(1, sizeof(struct TOPO_GROUP_NODE));
struct TOPO_NODE* ni = u_calloc(1, sizeof(struct TOPO_NODE));
if (!nq || !ni) { u_free(nq); u_free(ni); return NULL; }
ni->group_ref_count = 1; ni->node_id = nid; ni->ver = 0;
memcpy(ni->public_key, pubkey, SC_PUBKEY_SIZE);
nq->ll.size = sizeof(struct TOPO_GROUP_NODE) - sizeof(struct ll_entry);
nq->node_id = nid;
struct TOPO_SOCKMETA4* sm = memory_pool_alloc(inst->topo_groups->v4_sock_meta_pool);
if (sm) { sm->id = 0; sm->config_type = CFG_SERVER_TYPE_PUBLIC; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
struct TOPO_ADDR4* addr = memory_pool_alloc(inst->topo_groups->v4_addr_pool);
if (addr) { addr->addr[0] = a; addr->addr[1] = b; addr->addr[2] = c; addr->addr[3] = d; addr->port = port; addr->type = TOPO_ADDR_INTERFACE; addr->socket_id = 0; addr->protocol = TOPO_PROTO_UDP; addr->next = ni->v4_addrs; ni->v4_addrs = addr; }
if (inst && inst->topo_groups) topo_node_registry_store(inst->topo_groups, ni);
return nq;
}
static int create_temp_configs(void) {
if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp fail\n"); return -1; }
int base = 42000+(getpid()%15000);
@ -115,6 +97,8 @@ static int create_temp_configs(void) {
fprintf(f,"[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.0.2/24\ntun_ifname=tun98\n\n"
"[server:s_a]\naddr=127.0.0.1:%d\ntype=public\n\n[server:s_c]\naddr=127.0.0.1:%d\ntype=public\n",
privhex_b,pubhex_b,b_port_a,b_port_c);
fprintf(f, "[client:to_a]\npeer_public_key=%s\nlink=s_a:127.0.0.1:%d\n"
"[client:to_c]\npeer_public_key=%s\nlink=s_c:127.0.0.1:%d\n", pubhex_a, a_port, pubhex_c, c_port);
fclose(f);
snprintf(c_conf,sizeof(c_conf),"%s/c.conf",temp_dir);
@ -154,11 +138,9 @@ static int send_one_pkt(void) {
return ret;
}
static void connect_bgp_ready_cb(void* arg, struct ETCP_CONN* conn, int type) {
(void)conn;
if (type == ETCP_CONNECT_BGP_READY) {
int id = *(int*)arg; g_bgp_ready_mask |= (1 << id);
}
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);
}
static const char* state_name(int s) {
@ -179,7 +161,7 @@ static int all_acked(void) {
static void state_step(void) {
switch (g_state) {
case ST_INIT:
if (g_bgp_ready_mask == 3) {
if (transit_routes_ready()) {
struct ETCP_CONN* rc = topo_group_find_conn_for_node(topo_groups_get_default(g_a->topo_groups), node_c);
if (rc) g_state=ST_PHASE1_SEND;
}
@ -200,7 +182,7 @@ static void state_step(void) {
if (e) etcp_connection_close(((struct conn_queue_entry*)e->data)->conn);
e = queue_find_data_by_index(g_c->connections, (const uint8_t*)&node_b);
if (e) etcp_connection_close(((struct conn_queue_entry*)e->data)->conn); }
g_phase_sent=0; g_bgp_ready_mask=0; g_state=ST_B_RESTART;
g_phase_sent=0; g_state=ST_B_RESTART;
break;
case ST_B_RESTART:
#ifdef ROUTER_TEST_RECOVERY_ONLY
@ -219,14 +201,8 @@ static void state_step(void) {
g_b=utun_instance_create(ua,b_conf);
if (!g_b||utun_instance_init(g_b)<0) { fprintf(stderr,"FAIL: b restart\n"); g_test_ok=-1; g_fail_code=20; return; }
g_b->etcp_connect_timeout_tb = 300000;
struct TOPO_GROUP_NODE* nq_a = mknode(g_b, node_a, keys_a.public_key, 127,0,0,1, (uint16_t)a_port);
struct TOPO_GROUP_NODE* nq_c = mknode(g_b, node_c, keys_c.public_key, 127,0,0,1, (uint16_t)c_port);
if (!nq_a || !nq_c) { fprintf(stderr,"FAIL: mknode\n"); g_test_ok=-1; return; }
static int cc_a=0,cc_c=1;
etcp_connect(g_b,nq_a,connect_bgp_ready_cb,&cc_a,ETCP_CONNECT_BGP_READY);
etcp_connect(g_b,nq_c,connect_bgp_ready_cb,&cc_c,ETCP_CONNECT_BGP_READY);
}
if (g_bgp_ready_mask == 3) {
if (transit_routes_ready()) {
struct ETCP_CONN* rc = topo_group_find_conn_for_node(topo_groups_get_default(g_a->topo_groups), node_c);
if (rc) { g_phase_sent=0; g_state=ST_PHASE3_SEND; }
}
@ -278,14 +254,6 @@ int main(void) {
etcp_router_bind(g_c, TEST_SVC_ID, recv_handler);
struct TOPO_GROUP_NODE* nq_a = mknode(g_b, node_a, keys_a.public_key, 127,0,0,1, (uint16_t)a_port);
struct TOPO_GROUP_NODE* nq_c = mknode(g_b, node_c, keys_c.public_key, 127,0,0,1, (uint16_t)c_port);
if (!nq_a || !nq_c) { fprintf(stderr,"FAIL: mknode\n"); goto fail; }
g_b->etcp_connect_timeout_tb = 300000;
{ static int cc_a=0,cc_c=1;
etcp_connect(g_b,nq_a,connect_bgp_ready_cb,&cc_a,ETCP_CONNECT_BGP_READY);
etcp_connect(g_b,nq_c,connect_bgp_ready_cb,&cc_c,ETCP_CONNECT_BGP_READY); }
void* to_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS*10, NULL, timeout_cb, "to");
while (!g_test_ok) {
uasync_poll(ua, POLL_TB);

27
tests/test_group_ownership.c

@ -6,6 +6,26 @@
#include "transport_layer/node_conn_direct.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
static void check_deliberate_leave(struct TOPO_GROUP* group, struct ETCP_CONN* conn, enum topo_group_remove_reason reason) {
assert(topo_group_new_conn(group, conn) == 0);
struct SC_MYKEYS keys; assert(sc_generate_keypair(&keys) == SC_OK);
struct TOPO_NODE* node = u_calloc(1, sizeof(*node)); assert(node);
memcpy(node->public_key, keys.public_key, SC_PUBKEY_SIZE);
uint64_t id = node->node_id = sc_derive_node_id_from_pubkey(node->public_key);
assert(topo_node_registry_store(group->instance->topo_groups, node));
struct ll_entry* entry = queue_entry_new(sizeof(struct TOPO_GROUP_NODE) - sizeof(struct ll_entry)); assert(entry);
struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)entry;
nq->node_id = id;
uint64_t hops[] = { id, conn->peer_node_id };
assert(topo_group_add_path(nq, conn, hops, 2, 100) == 0);
assert(queue_data_put_with_index(group->nodes, entry) == 0);
topo_group_remove_conn(group, conn, reason);
assert(!topo_node_find_by_id(group, id));
assert(!group->recovery_list && !topo_node_registry_find(group->instance->topo_groups, id));
assert(!conn->close_requested);
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO);
@ -28,6 +48,13 @@ int main(void) {
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);
/* Создание группы при уже существующем транспорте читает payload записи conn. */
struct TOPO_GROUP* c = topo_groups_create_group(inst->topo_groups, 44, TOPO_GROUP_TYPE_UTUN, NULL); assert(c);
assert(queue_entry_count(c->senders_list) == 0);
check_deliberate_leave(c, conn, TOPO_REMOVE_LOCAL_LEAVE);
check_deliberate_leave(c, conn, TOPO_REMOVE_REMOTE_LEAVE);
check_deliberate_leave(c, conn, TOPO_REMOVE_MEMBER_INVALID);
topo_groups_remove_group(inst->topo_groups, 44);
/* Проверяем владение независимо от доставки пакетов и ответа удалённого узла. */
assert(topo_group_new_conn(a, conn) == 0);
assert(topo_group_new_conn(a, conn) == 0 && queue_entry_count(a->senders_list) == 1);

96
tests/test_ncd_config.c

@ -0,0 +1,96 @@
#include <assert.h>
#include <stdio.h>
#include <string.h>
#include "utun_instance.h"
#include "config_parser.h"
#include "transport_layer/etcp.h"
#include "transport_layer/etcp_connections.h"
#include "transport_layer/node_conn_direct.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
static int links(struct ETCP_CONN* conn) {
int count = 0;
for (struct ETCP_LINK* link = conn->links; link; link = link->next) count++;
return count;
}
static void endpoint(struct CFG_CLIENT_LINK* link, struct CFG_SERVER* local, int port) {
memset(link, 0, sizeof(*link)); link->local_srv = local;
struct sockaddr_in* addr = (struct sockaddr_in*)&link->remote_addr;
addr->sin_family = AF_INET; addr->sin_addr.s_addr = htonl(0x7f000001); addr->sin_port = htons((uint16_t)port);
}
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"
"[server: tcp]\naddr=127.0.0.1:0\ntype=public\ntransport=tcp\n");
assert(inst && utun_instance_init(inst) >= 0);
struct CFG_SERVER *udp = NULL, *tcp = NULL;
for (struct CFG_SERVER* s = inst->config->servers; s; s = s->next) {
if (!strcmp(s->name, "udp")) udp = s;
if (!strcmp(s->name, "tcp")) tcp = s;
}
assert(udp && tcp);
struct CFG_CLIENT client = { .name = "peer" };
struct SC_MYKEYS peer; assert(sc_generate_keypair(&peer) == SC_OK);
for (int i = 0; i < 32; i++) snprintf(client.peer_public_key_hex + i * 2, 3, "%02x", peer.public_key[i]);
struct CFG_CLIENT_LINK first, second, third;
endpoint(&first, udp, 9); endpoint(&second, udp, 10); endpoint(&third, tcp, 11);
first.next = &second; second.next = &third; client.links = &first;
assert(init_client_connections(inst, &client) == 0);
assert(inst->config_conn_handles && !inst->config_conn_handles->next);
struct ETCP_CONN* conn = node_conn_direct_get_conn(inst->config_conn_handles->handle);
assert(conn && links(conn) == 3 && !conn->fin_wait);
struct NODE_CONN_DIRECT* chat = NULL;
assert(node_conn_direct_open(inst, conn->peer_node_id, NULL, NULL, &chat, NULL) == NCD_REUSED);
assert(init_client_connections(inst, &client) == 0);
assert(node_conn_direct_get_conn(inst->config_conn_handles->handle) == conn && links(conn) == 3 && !conn->fin_wait);
/* Ошибка нового набора не освобождает предыдущие handles. */
struct CFG_CLIENT invalid = { .name = "invalid", .peer_public_key_hex = "not-a-key" };
client.next = &invalid;
struct CONFIG_CONN_HANDLE* previous = inst->config_conn_handles;
assert(init_client_connections(inst, &client) < 0 && inst->config_conn_handles == previous);
assert(node_conn_direct_get_conn(chat) == conn && !conn->fin_wait);
client.next = NULL;
/* REALITY идёт по тому же NCD пути; полное SNI не обрезается. */
client.links = &third; client.reality_enabled = 1;
reality_client_config_set_defaults(&client.reality);
memcpy(client.reality.server_static_pubkey, peer.public_key, 32);
memset(client.reality.server_name, 'a', 180); client.reality.server_name[180] = 0;
assert(init_client_connections(inst, &client) == 0);
assert(links(conn) == 3);
struct ETCP_LINK* real = NULL;
for (struct ETCP_LINK* link = conn->links; link; link = link->next) if (link->is_tcp) real = link;
assert(real && real->reality_set && !strcmp(real->reality.server_name, client.reality.server_name));
assert(init_client_connections(inst, &client) == 0 && links(conn) == 3);
assert(init_client_connections(inst, NULL) == 0 && !inst->config_conn_handles);
assert(node_conn_direct_get_conn(chat) == conn && !conn->fin_wait);
node_conn_direct_force_close(chat);
uasync_poll(ua, 0);
/* Клиентский TCP/REALITY не требует локального TCP listen-сокета. */
struct ETCP_SOCKET* sock = inst->etcp_sockets;
while (sock) {
struct ETCP_SOCKET* next = sock->next;
if (sock->is_tcp) etcp_socket_remove(sock);
sock = next;
}
third.local_srv = NULL;
assert(init_client_connections(inst, &client) == 0);
conn = node_conn_direct_get_conn(inst->config_conn_handles->handle);
assert(conn && links(conn) == 1 && conn->links->is_tcp && !conn->links->conn);
assert(conn->links->reality_set && !strcmp(conn->links->reality.server_name, client.reality.server_name));
assert(init_client_connections(inst, NULL) == 0);
inst->running = 0; utun_instance_destroy(inst);
uasync_poll(ua, 0);
uasync_destroy(ua, 0);
return 0;
}

7
tests/test_reality_bgp.c

@ -38,7 +38,9 @@ int main(void) {
r1.short_id[0] = 0xAA; r2.short_id[0] = 0xBB;
r1.version[0] = 1; r1.version[1] = 2; r1.version[2] = 3;
for (int i = 0; i < 32; i++) r1.server_pubkey[i] = (uint8_t)(i + 1);
strcpy(r1.server_name, "www.example.com");
char long_sni[181]; memset(long_sni, 'a', sizeof(long_sni) - 1); long_sni[180] = 0;
long_sni[50] = '.'; long_sni[100] = '.'; long_sni[150] = '.';
strcpy(r1.server_name, long_sni);
strcpy(r2.server_name, "api.example.com");
r1.next = &r2; r2.next = NULL;
ni.reality_socks = &r1;
@ -56,7 +58,7 @@ int main(void) {
CHECK(memcmp(link.reality.short_id, r1.short_id, 8) == 0, "apply short_id");
CHECK(memcmp(link.reality.server_static_pubkey, r1.server_pubkey, 32) == 0, "apply pubkey");
CHECK(memcmp(link.reality.version, r1.version, 3) == 0, "apply version");
CHECK(strcmp(link.reality.server_name, "www.example.com") == 0, "apply server_name");
CHECK(strcmp(link.reality.server_name, long_sni) == 0, "apply full server_name");
topo_node_apply_reality(&link, NULL);
CHECK(link.reality_set == 1, "apply NULL rs is no-op");
@ -82,6 +84,7 @@ int main(void) {
}
const struct TOPO_REALITY_SOCK* rr1 = topo_node_find_reality_sock(out_ni, 3);
CHECK(rr1 != NULL && memcmp(rr1->server_pubkey, r1.server_pubkey, 32) == 0, "round-trip: pubkey");
CHECK(rr1 != NULL && !strcmp(rr1->server_name, long_sni), "round-trip: full server_name");
CHECK(rr1 != NULL && rr1->version[0] == 1 && rr1->version[1] == 2 && rr1->version[2] == 3, "round-trip: version");
}

Loading…
Cancel
Save