Browse Source

Apply fresh BGP addresses to shared NCD connections

proxy
evgeny 5 days ago
parent
commit
26e6c05313
  1. 11
      doc/node_snapshot.md
  2. 1
      src/routing_layer/topo_group.c
  3. 123
      src/transport_layer/node_conn_direct.c
  4. 4
      src/transport_layer/node_conn_direct.h
  5. 14
      tests/test_etcp_stcp.c
  6. 42
      tests/test_node_conn_direct.c

11
doc/node_snapshot.md

@ -23,3 +23,14 @@ Timestamp 0 обозначает bootstrap-сведения без подпис
начальные адреса). Они могут помочь установить первое соединение, но не
заменяют принятую подписанную запись. По BGP принимаются подписанные записи
с ненулевым timestamp. Формат NODEINFO изменён без обратной совместимости.
После принятия BGP-записи `node_conn_direct_update_node` применяет адреса к
существующему NCD-соединению. Отсутствующее соединение этот вызов не создаёт.
Линки дополняются как при UP, так и до установления соединения; повторная
доставка не создаёт дубликаты UDP/TCP. Работающие линки не закрываются только
из-за смены списка адресов. NCD удерживает ссылку на общую запись, пока
обслуживает соединение, и использует её при восстановлении.
`node_conn_direct_open_node` сверяет переданную запись с реестром/БД. Более
свежая подписанная запись принимается целиком, старая уступает уже известной.
Bootstrap-адреса не перекрывают подписанную запись даже при DOWN.

1
src/routing_layer/topo_group.c

@ -1008,6 +1008,7 @@ 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;

123
src/transport_layer/node_conn_direct.c

@ -41,6 +41,8 @@ struct ncd_entry {
struct UTUN_INSTANCE* inst;
unsigned dispatch_depth;
uint8_t closed;
uint8_t node_ref;
uint64_t node_timestamp;
int handle_count;
struct NODE_CONN_DIRECT* handles; /* связный список всех handle'ов */
@ -154,6 +156,7 @@ static int ncd_add_udp_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* use_sock
(unsigned long long)conn->peer_node_id, *(const uint64_t*)conn->crypto_ctx.peer_public_key);
etcp_link_close(stale);
etcp_cbk_fire(old, ETCP_CBK_EVENT_NODE_CHANGED);
stale = NULL;
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "[ncd] addr conflict: port already used by live conn 0x%016llx for node 0x%016llx",
(unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id);
@ -174,10 +177,50 @@ static int ncd_add_udp_link(struct ETCP_CONN* conn, struct ETCP_SOCKET* use_sock
return etcp_link_new(conn, use_sock, sa, 0) ? 1 : 0;
}
static int ncd_same_address(const struct sockaddr_storage* a, const struct sockaddr_storage* b) {
if (a->ss_family != b->ss_family) return 0;
if (a->ss_family == AF_INET) {
const struct sockaddr_in* x = (const struct sockaddr_in*)a;
const struct sockaddr_in* y = (const struct sockaddr_in*)b;
return x->sin_port == y->sin_port && x->sin_addr.s_addr == y->sin_addr.s_addr;
}
if (a->ss_family == AF_INET6) {
const struct sockaddr_in6* x = (const struct sockaddr_in6*)a;
const struct sockaddr_in6* y = (const struct sockaddr_in6*)b;
return x->sin6_port == y->sin6_port && x->sin6_scope_id == y->sin6_scope_id &&
!memcmp(&x->sin6_addr, &y->sin6_addr, sizeof(x->sin6_addr));
}
return 0;
}
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;
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);
etcp_tcp_link_start_connect(link, addr, 0);
return 1;
}
static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
struct ETCP_SOCKET* specific_sock) {
if (!ni || ni->timestamp < entry->node_timestamp) {
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] ignore stale addresses node=%016llx current=%llu incoming=%llu",
(unsigned long long)entry->node_id, (unsigned long long)entry->node_timestamp,
(unsigned long long)(ni ? ni->timestamp : 0)); return 0;
}
struct ETCP_CONN* conn = entry->conn;
struct UTUN_INSTANCE* inst = conn->instance;
if (memcmp(ni->public_key, conn->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] address update key mismatch node=%016llx", (unsigned long long)entry->node_id);
return -1;
}
if (!entry->node_ref && topo_node_registry_find(inst->topo_groups, entry->node_id) == ni) {
topo_node_registry_ref(inst->topo_groups, entry->node_id); entry->node_ref = 1;
}
entry->node_timestamp = ni->timestamp;
int link_count = 0, v4_skip = 0, v6_skip = 0, v4_addrs = 0, v6_addrs = 0;
uint64_t nid = entry->node_id;
@ -202,9 +245,9 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
if (a->protocol & TOPO_PROTO_TCP) {
for (int vi = 0; vi < nviews; vi++) {
const struct sock_view* v = &views[vi];
if (!v->is_tcp || !sock_match(v, a->type, tcfg, tnat, tip)) continue;
struct ETCP_LINK* tlink = etcp_link_new(conn, v->sock, &sa, 0);
if (tlink) { tlink->is_tcp = 1; topo_node_apply_reality(tlink, topo_node_find_reality_sock(ni, a->socket_id)); etcp_tcp_link_start_connect(tlink, &sa, a->port); link_count++; }
if (!v->is_tcp || (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 (a->protocol & TOPO_PROTO_UDP) {
@ -235,9 +278,11 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
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 || s->local_addr.ss_family != AF_INET6 || s->type == CFG_SERVER_TYPE_PRIVATE) continue;
struct ETCP_LINK* tlink = etcp_link_new(conn, s, &sa, 0);
if (tlink) { tlink->is_tcp = 1; topo_node_apply_reality(tlink, topo_node_find_reality_sock(ni, a->socket_id)); etcp_tcp_link_start_connect(tlink, &sa, a->port); link_count++; }
if (!s->is_tcp || (specific_sock && s != specific_sock) || s->local_addr.ss_family != AF_INET6 ||
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 (a->protocol & TOPO_PROTO_UDP) {
@ -256,13 +301,24 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni,
}
}
if (link_count == 0) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] create_links: node=0x%016llx v4=%d/skp=%d v6=%d/skp=%d → %d links (UDP+TCP)",
if (link_count > 0) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] create_links: node=0x%016llx v4=%d/skp=%d v6=%d/skp=%d → %d new links (UDP+TCP)",
(unsigned long long)nid, v4_addrs, v4_skip, v6_addrs, v6_skip, link_count);
}
return link_count;
}
int node_conn_direct_update_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] update without instance"); return -1; }
struct ncd_entry* entry = ncd_registry_find(inst, node_id);
struct TOPO_NODE* ni = topo_node_registry_find(inst->topo_groups, node_id);
if (!entry || !ni) return 0;
int added = ncd_create_links(entry, ni, NULL);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] apply node=%016llx timestamp=%llu up=%d added=%d",
(unsigned long long)node_id, (unsigned long long)ni->timestamp, entry->up, added);
return added;
}
/* ═══════════ единая диспетчеризация событий ═══════════ */
/*
* Центральный диспетчер: принимает событие (UP/DOWN/TIMEOUT), проверяет что оно
@ -355,6 +411,9 @@ static void ncd_finish(struct ncd_entry* entry, int close_conn) {
if (entry->closed) return;
entry->closed = 1;
ncd_registry_remove(entry->inst, entry);
if (entry->node_ref) {
topo_node_registry_unref(entry->inst->topo_groups, entry->node_id); entry->node_ref = 0;
}
if (entry->connect_timer) { uasync_cancel_timeout(entry->ua, entry->connect_timer); entry->connect_timer = NULL; }
if (entry->fin_wait_timer) { uasync_cancel_timeout(entry->ua, entry->fin_wait_timer); entry->fin_wait_timer = NULL; }
struct ETCP_CONN* conn = entry->conn; entry->conn = NULL;
@ -638,7 +697,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
/* 2. Ищем conn через instance_find_conn (входящее / созданное etcp_connect) */
struct ETCP_CONN* conn = instance_find_conn(inst, node_id);
if (conn && conn->state == 1 && !conn->close_requested) {
if (conn && conn->state != 2 && !conn->close_requested) {
/* снять fin_wait если был */
if (conn->fin_wait) {
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open REUSED clearing fin_wait (existing conn) node=0x%016llx", (unsigned long long)node_id);
@ -760,7 +819,7 @@ int node_conn_direct_open(struct UTUN_INSTANCE* inst, uint64_t node_id,
* известны заранее (например из BGP-анонса). ni не владеем — можно
* передать временную структуру, копия не делается.
*/
int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
static int ncd_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
ncd_callback cb, void* cb_arg,
struct NODE_CONN_DIRECT** out_handle,
struct TOPO_NODE* ni,
@ -784,12 +843,14 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
entry->handle_count++;
*out_handle = h;
if (entry->conn && entry->conn->state == 1 && entry->up) {
ncd_create_links(entry, ni, specific_sock);
ncd_schedule_deliver_up(inst->ua, h);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED (ready) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
} else if (entry->timed_out) {
ncd_revive_entry(entry, inst, ni, specific_sock);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED (revived) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
} else {
ncd_create_links(entry, ni, specific_sock);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED (pending) node=0x%016llx handles=%d", (unsigned long long)node_id, entry->handle_count);
}
return NCD_REUSED;
@ -821,6 +882,7 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
etcp_conn_add_cbk(conn, ncd_down_cb, entry, ETCP_CBK_EVENT_DOWN);
if (conn->state == 1 && entry->up) {
ncd_create_links(entry, ni, specific_sock);
ncd_schedule_deliver_up(inst->ua, h);
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node REUSED new-entry (ready) node=0x%016llx conn=%p", (unsigned long long)node_id, (void*)conn);
} else {
@ -888,7 +950,7 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
int link_count = ncd_create_links(entry, ni, specific_sock);
if (link_count == 0)
DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] open_node (UDP only, TCP added by conn_mgr) node=0x%016llx links=%d", (unsigned long long)node_id, link_count);
DEBUG_DEBUG(DEBUG_CATEGORY_NCD, "[ncd] open_node: no new compatible links node=0x%016llx", (unsigned long long)node_id);
DEBUG_DEBUG(DEBUG_CATEGORY_TIMERS, "[ncd] connect timer SET: node=0x%016llx value_tb=%u now_tb=%llu path=NEW_node_conn_direct_open_node",
(unsigned long long)node_id, inst->etcp_connect_timeout_tb, (unsigned long long)get_time_tb());
@ -899,6 +961,45 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
}
}
int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
ncd_callback cb, void* cb_arg, struct NODE_CONN_DIRECT** out_handle,
struct TOPO_NODE* ni, struct ETCP_SOCKET* specific_sock) {
if (out_handle) *out_handle = NULL;
if (!inst || !out_handle || !ni || node_id != ni->node_id ||
node_id != sc_derive_node_id_from_pubkey(ni->public_key) || (ni->timestamp && topo_node_verify(ni) < 0)) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] invalid node passed to open"); return NCD_ERR;
}
struct TOPO_NODE* known = ncd_lookup_node(inst, node_id);
if (ni->timestamp && (!known || ni->timestamp > known->timestamp)) {
uint8_t wire[4096];
struct TOPO_GROUP_NODE nq = {0};
struct TOPO_GROUP group = { .instance = inst };
struct TOPO_NODE* copy = NULL;
int len = topo_node_serialize(ni, &nq, 0, 0, wire, sizeof(wire), 0);
if (len < 0 || !inst->topo_groups ||
topo_node_deserialize(&group, wire, (size_t)len, &copy, NULL, NULL, NULL, NULL) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] cannot copy node snapshot");
if (known) topo_node_registry_unref(inst->topo_groups, node_id);
return NCD_ERR;
}
struct TOPO_NODE* stored = topo_node_registry_store(inst->topo_groups, copy);
if (known) topo_node_registry_unref(inst->topo_groups, node_id);
known = stored;
if (!stored) { topo_node_destroy(inst->topo_groups, copy); return NCD_ERR; }
if (inst->topo_sqlite_db && topo_node_sqlite_snapshot_put(inst->topo_sqlite_db, stored) < 0)
DEBUG_ERROR(DEBUG_CATEGORY_NCD, "[ncd] cannot persist received node snapshot");
}
if (known && ni->timestamp && known->timestamp == ni->timestamp &&
memcmp(known->x25519_self_sig, ni->x25519_self_sig, 64)) {
DEBUG_WARN(DEBUG_CATEGORY_NCD, "[ncd] conflicting node snapshot node=%016llx", (unsigned long long)node_id);
topo_node_registry_unref(inst->topo_groups, node_id); return NCD_ERR;
}
struct TOPO_NODE* selected = known && known->timestamp ? known : ni;
int result = ncd_open_node(inst, node_id, cb, cb_arg, out_handle, selected, specific_sock);
if (known) topo_node_registry_unref(inst->topo_groups, node_id);
return result;
}
/*
* Закрыть handle (graceful shutdown).
*

4
src/transport_layer/node_conn_direct.h

@ -67,6 +67,10 @@ int node_conn_direct_open_node(struct UTUN_INSTANCE* inst, uint64_t node_id,
void node_conn_direct_close(struct NODE_CONN_DIRECT* h);
/* Применить принятую запись из общего реестра к уже существующему соединению.
* Не создаёт соединение и не приобретает handle. Возвращает число новых линков. */
int node_conn_direct_update_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Force close без fin_wait: немедленно удаляет ETCP-коллбэки и, если это
* последний handle, закрывает соединение и освобождает ncd_entry.
* Используется при уничтожении CM entry — гарантирует, что после вызова

14
tests/test_etcp_stcp.c

@ -6,6 +6,7 @@
#include "secure_channel.h"
#include "topo_group.h"
#include "topo_node.h"
#include "node_conn_direct.h"
#include "../src/utun_instance.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
@ -119,6 +120,18 @@ int main(void) {
int ticks = 0;
while (!cli_conn && ticks < 5000) { uasync_poll(ua, 10); ticks++; }
TASSERT(cli_conn);
int initial_links = 0;
for (struct ETCP_LINK* link = cli_conn->links; link; link = link->next) initial_links++;
struct NODE_CONN_DIRECT *owner1 = NULL, *owner2 = NULL;
struct TOPO_NODE* peer = topo_node_registry_find(cli_inst->topo_groups, srv_node->node_id);
TASSERT(node_conn_direct_open_node(cli_inst, peer->node_id, NULL, NULL, &owner1, peer, NULL) == NCD_REUSED);
TASSERT(node_conn_direct_open_node(cli_inst, peer->node_id, NULL, NULL, &owner2, peer, NULL) == NCD_REUSED);
TASSERT(node_conn_direct_update_node(cli_inst, peer->node_id) == 0);
int repeated_links = 0;
for (struct ETCP_LINK* link = cli_conn->links; link; link = link->next) repeated_links++;
TASSERT(repeated_links == initial_links);
node_conn_direct_close(owner1);
TASSERT(node_conn_direct_get_conn(owner2) == cli_conn);
{ struct ll_entry* e = srv_inst->connections->head;
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
@ -140,6 +153,7 @@ int main(void) {
while (g_srv_recv_count < 1 && ticks < 2000) { uasync_poll(ua, 10); ticks++; }
TASSERT(g_srv_recv_count == 1);
TASSERT(memcmp(g_srv_recv_data, msg, strlen(msg)) == 0);
node_conn_direct_close(owner2);
srv_inst->running = 0; cli_inst->running = 0;
utun_instance_destroy(srv_inst); srv_inst = NULL;

42
tests/test_node_conn_direct.c

@ -95,6 +95,46 @@ static void t6(void* arg); static void t7(void* arg);
static void t8(void* arg); static void t9(void* arg);
static void t_done(void* arg);
/* Новая подписанная запись должна добавлять адрес и на pending, и на UP conn. */
static int update_addresses(int up) {
struct ETCP_CONN* conn = node_conn_direct_get_conn(gh[0]);
int before = 0;
for (struct ETCP_LINK* l = conn->links; l; l = l->next) before++;
struct TOPO_NODE* source = topo_node_registry_find(g_b->topo_groups, nid_b);
struct TOPO_GROUP_NODE nq = {0};
struct TOPO_GROUP group = { .instance = g_a, .group_type = TOPO_GROUP_TYPE_CHAT, .group_id = 42 };
uint8_t wire[4096];
int len = topo_node_serialize(source, &nq, 42, 0, wire + 2, sizeof(wire) - 2, 0);
struct TOPO_NODE* copy = NULL;
if (len < 0 || topo_node_deserialize(&group, wire + 2, (size_t)len, &copy, NULL, NULL, NULL, NULL) < 0) return -1;
struct TOPO_ADDR4* addr = memory_pool_alloc(g_a->topo_groups->v4_addr_pool);
if (!addr) { topo_node_destroy(g_a->topo_groups, copy); return -1; }
*addr = (struct TOPO_ADDR4){ .next = copy->v4_addrs, .addr = {127,0,0,2}, .port = (uint16_t)(pb + up),
.protocol = TOPO_PROTO_UDP };
copy->v4_addrs = addr;
copy->timestamp += (uint64_t)up + 1;
uint8_t message[TOPO_SIG_MSG_MAX_SIZE];
len = topo_node_build_sig_msg(copy, message, sizeof(message));
if (len < 0 || sc_ed25519_sign(g_b->my_ed25519_privkey, message, (size_t)len, copy->x25519_self_sig) != SC_OK) return -1;
len = topo_node_serialize(copy, &nq, 42, 0, wire + 2, sizeof(wire) - 2, 0);
group.nodes = queue_new(ua, 16, 0, 8, "ncd_address_test");
int result = topo_group_process_nodeinfo(&group, conn, wire, (size_t)len + 2);
int after = 0;
for (struct ETCP_LINK* l = conn->links; l; l = l->next) after++;
if (result != 0 || after != before + 1 || node_conn_direct_update_node(g_a, nid_b) != 0) result = -1;
struct NODE_CONN_DIRECT* temporary = NULL;
if (node_conn_direct_open_node(g_a, nid_b, NULL, NULL, &temporary, source, NULL) != NCD_REUSED) result = -1;
if (temporary) node_conn_direct_close(temporary);
int repeated = 0;
for (struct ETCP_LINK* l = conn->links; l; l = l->next) repeated++;
if (repeated != after || (up && !conn->links_up)) result = -1;
struct TOPO_GROUP_NODE* peer = topo_node_find_by_id(&group, nid_b);
if (peer) { topo_group_remove_path(peer, conn); topo_nodeq_remove_node(&group, peer); }
queue_free(group.nodes);
topo_node_destroy(g_a->topo_groups, copy);
return result;
}
static void t1(void* arg) {
(void)arg; g_phase = 1;
fprintf(stderr, "\n=== P1: two opens before poll ===\n"); fflush(stderr);
@ -108,6 +148,7 @@ static void t1(void* arg) {
fprintf(stderr, " open#1 → %s\n", r == NCD_NEW ? "NCD_NEW" : r == NCD_REUSED ? "NCD_REUSED" : "ERR");
if (r != NCD_REUSED) { fail("expected NCD_REUSED"); return; }
if (update_addresses(0) != 0) { fail("pending address update"); return; }
uasync_call_soon(ua, NULL, t2);
}
@ -125,6 +166,7 @@ static void t2(void* arg) {
static void t3(void* arg) {
(void)arg; g_phase = 3;
fprintf(stderr, "\n=== P3: open on ready conn ===\n"); fflush(stderr);
if (update_addresses(1) != 0) { fail("UP address update"); return; }
int r = node_conn_direct_open(g_a, nid_b, ncd_cb, (void*)2, &gh[2], NULL);
fprintf(stderr, " open#2 → %s\n", r == NCD_REUSED ? "NCD_REUSED" : "?");
if (r != NCD_REUSED) { fail("expected NCD_REUSED"); return; }

Loading…
Cancel
Save