Browse Source

topo: держать актуальные self-адреса в node_addresses, обновлять/рассылать только при изменении

topo_node_update_my_addresses теперь собирает адреса во временные списки,
сравнивает со старыми и при отсутствии изменений выходит без бампа ver,
записи в БД и broadcast. Гейт member_exists убран: self-адреса пишутся
в node_addresses при любом запуске/изменении (с гарантией строки nodes для FK).
addrs_put логирует ошибки INSERT.
topo_upd
evgeny 2 months ago
parent
commit
eb24736e27
  1. 117
      src/routing_layer/topo_node.c
  2. 2
      src/routing_layer/topo_node.h
  3. 10
      src/routing_layer/topo_node_sqlite.c

117
src/routing_layer/topo_node.c

@ -16,6 +16,7 @@
#include "topo_node_sqlite.h"
#include "../transport_layer/etcp_connections.h"
#include "../lib/u_async.h"
#include "../ntp_time.h"
void topo_node_ref(struct TOPO_NODE* ni) {
if (!ni) return;
@ -54,6 +55,39 @@ static void free_v6_sub_list(struct memory_pool* pool, struct TOPO_SUBNET6* head
while (head) { struct TOPO_SUBNET6* next = head->next; memory_pool_free(pool, head); head = next; }
}
/* Упорядоченное сравнение списков адресов/мета (порядок обхода сокетов детерминирован,
prepend даёт одинаковый порядок при неизменном наборе сокетов). */
static int sockmeta4_list_equal(struct TOPO_SOCKMETA4* a, struct TOPO_SOCKMETA4* b) {
while (a && b) {
if (a->id != b->id || a->config_type != b->config_type || a->nat_type != b->nat_type) return 0;
a = a->next; b = b->next;
}
return a == NULL && b == NULL;
}
static int addr4_list_equal(struct TOPO_ADDR4* a, struct TOPO_ADDR4* b) {
while (a && b) {
if (a->port != b->port || a->type != b->type || a->socket_id != b->socket_id ||
a->protocol != b->protocol || memcmp(a->addr, b->addr, 4) != 0) return 0;
a = a->next; b = b->next;
}
return a == NULL && b == NULL;
}
static int sockmeta6_list_equal(struct TOPO_SOCKMETA6* a, struct TOPO_SOCKMETA6* b) {
while (a && b) {
if (a->id != b->id || a->config_type != b->config_type || a->nat_type != b->nat_type) return 0;
a = a->next; b = b->next;
}
return a == NULL && b == NULL;
}
static int addr6_list_equal(struct TOPO_ADDR6* a, struct TOPO_ADDR6* b) {
while (a && b) {
if (a->port != b->port || a->type != b->type || a->socket_id != b->socket_id ||
a->protocol != b->protocol || memcmp(a->addr, b->addr, 16) != 0) return 0;
a = a->next; b = b->next;
}
return a == NULL && b == NULL;
}
/* извлекает hop_list из лучшего живого path в nq->paths (или любого, если живых нет).
память принадлежит TOPO_NODEPATH в paths, не освобождать */
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt) {
@ -681,17 +715,6 @@ int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size) {
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group) {
if (!instance || !group) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; }
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: ENTER inst=%p nid=0x%016llx grp=%p grp_inst=%p grp_type=%d",
(void*)instance, (unsigned long long)instance->node_id, (void*)group, (void*)group->instance, group->group_type);
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: topo_groups=%p etcp=%p reg=%p",
(void*)instance->topo_groups, (void*)instance->etcp_sockets,
instance->topo_groups ? (void*)instance->topo_groups->node_registry : NULL);
if (instance->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: pools sm4=%p a4=%p sm6=%p a6=%p sub4=%p sub6=%p",
(void*)instance->topo_groups->v4_sock_meta_pool, (void*)instance->topo_groups->v4_addr_pool,
(void*)instance->topo_groups->v6_sock_meta_pool, (void*)instance->topo_groups->v6_addr_pool,
(void*)instance->topo_groups->v4_subnet_pool, (void*)instance->topo_groups->v6_subnet_pool);
}
size_t name_len = 0;
if (instance->name[0]) { name_len = strlen(instance->name); if (name_len > 63) name_len = 63; }
int vc = 0, vc6 = 0;
@ -790,11 +813,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
lq->last_ver = saved_ver; }
e_sock = instance->etcp_sockets;
int etcp_iter_cnt = 0;
while (e_sock) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: etcp_iter[%d] e_sock=%p type=%d fam=%d next=%p",
etcp_iter_cnt, (void*)e_sock, e_sock->type, e_sock->local_addr.ss_family, (void*)e_sock->next);
etcp_iter_cnt++;
if (e_sock->is_tcp) { e_sock = e_sock->next; continue; }
if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; }
if (e_sock->local_addr.ss_family == AF_INET) {
@ -831,10 +850,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
e_sock = e_sock->next;
}
{ struct ETCP_SOCKET* ts = instance->etcp_sockets;
int tcp_iter_cnt = 0;
while (ts) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "TOPO_UPD_SELF: tcp_iter[%d] ts=%p type=%d fam=%d next=%p",
tcp_iter_cnt, (void*)ts, ts->type, ts->local_addr.ss_family, (void*)ts->next);
tcp_iter_cnt++; if (!ts->is_tcp) { ts = ts->next; continue; }
while (ts) { if (!ts->is_tcp) { ts = ts->next; continue; }
if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; }
struct sockaddr_storage* addr = ts->interface_addr.ss_family ? &ts->interface_addr : &ts->local_addr;
if (!addr || !addr->ss_family) { ts = ts->next; continue; }
@ -895,19 +911,18 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
return vc;
}
void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; }
int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (!instance || !instance->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return -1; }
struct TOPO_GROUP* default_group = topo_groups_get_default(instance->topo_groups);
if (!default_group || !default_group->local_node) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "no local_node"); return; }
if (!default_group || !default_group->local_node) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "no local_node"); return -1; }
struct TOPO_NODE* ni = topo_node_registry_find(instance->topo_groups, instance->node_id);
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "my node not in registry"); return; }
if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "my node not in registry"); return -1; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "updating my addresses, old ver=%d", ni->ver);
free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, ni->v4_sock_meta); ni->v4_sock_meta = NULL;
free_v4_addr_list(instance->topo_groups->v4_addr_pool, ni->v4_addrs); ni->v4_addrs = NULL;
free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, ni->v6_sock_meta); ni->v6_sock_meta = NULL;
free_v6_addr_list(instance->topo_groups->v6_addr_pool, ni->v6_addrs); ni->v6_addrs = NULL;
/* Собираем новые списки во временные головы, чтобы сравнить со старыми до их освобождения. */
struct TOPO_SOCKMETA4* m4_head = NULL;
struct TOPO_ADDR4* v4_head = NULL;
struct TOPO_SOCKMETA6* m6_head = NULL;
struct TOPO_ADDR6* v6_head = NULL;
struct ETCP_SOCKET* e_sock = instance->etcp_sockets;
while (e_sock) {
@ -916,7 +931,7 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
if (e_sock->local_addr.ss_family == AF_INET) {
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(instance->topo_groups->v4_sock_meta_pool);
sm->id = e_sock->sock_id; sm->config_type = e_sock->type; sm->nat_type = e_sock->nat_type;
sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
sm->next = m4_head; m4_head = sm; }
struct sockaddr_in* local_sin = (struct sockaddr_in*)&e_sock->local_addr;
struct sockaddr_in* if_sin = (struct sockaddr_in*)&e_sock->interface_addr;
struct sockaddr_in* nat_sin = (struct sockaddr_in*)&e_sock->nat_addr;
@ -932,16 +947,16 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
a->type = TOPO_ADDR_INTERFACE; a->socket_id = e_sock->sock_id;
}
a->protocol = TOPO_PROTO_UDP;
a->next = ni->v4_addrs; ni->v4_addrs = a; }
a->next = v4_head; v4_head = a; }
} else if (e_sock->local_addr.ss_family == AF_INET6) {
{ struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(instance->topo_groups->v6_sock_meta_pool);
sm6->id = e_sock->sock_id; sm6->config_type = e_sock->type; sm6->nat_type = e_sock->nat_type;
sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; }
sm6->next = m6_head; m6_head = sm6; }
struct sockaddr_in6* if_sin6 = (struct sockaddr_in6*)&e_sock->interface_addr;
{ struct TOPO_ADDR6* a6 = memory_pool_alloc(instance->topo_groups->v6_addr_pool);
memcpy(a6->addr, &if_sin6->sin6_addr, 16); a6->port = ntohs(if_sin6->sin6_port);
a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = e_sock->sock_id; a6->protocol = TOPO_PROTO_UDP;
a6->next = ni->v6_addrs; ni->v6_addrs = a6; }
a6->next = v6_head; v6_head = a6; }
}
e_sock = e_sock->next;
}
@ -954,31 +969,52 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
struct sockaddr_in* sin = (struct sockaddr_in*)addr;
{ struct TOPO_SOCKMETA4* sm = memory_pool_alloc(instance->topo_groups->v4_sock_meta_pool);
sm->id = ts->sock_id; sm->config_type = ts->type; sm->nat_type = ts->nat_type;
sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; }
sm->next = m4_head; m4_head = sm; }
{ struct TOPO_ADDR4* a = memory_pool_alloc(instance->topo_groups->v4_addr_pool);
memcpy(a->addr, &sin->sin_addr.s_addr, 4); a->port = ntohs(sin->sin_port);
a->type = TOPO_ADDR_INTERFACE; a->socket_id = ts->sock_id; a->protocol = TOPO_PROTO_TCP;
a->next = ni->v4_addrs; ni->v4_addrs = a; }
a->next = v4_head; v4_head = a; }
} else if (addr->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)addr;
{ struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(instance->topo_groups->v6_sock_meta_pool);
sm6->id = ts->sock_id; sm6->config_type = ts->type; sm6->nat_type = ts->nat_type;
sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; }
sm6->next = m6_head; m6_head = sm6; }
{ struct TOPO_ADDR6* a6 = memory_pool_alloc(instance->topo_groups->v6_addr_pool);
memcpy(a6->addr, &sin6->sin6_addr, 16); a6->port = ntohs(sin6->sin6_port);
a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = ts->sock_id; a6->protocol = TOPO_PROTO_TCP;
a6->next = ni->v6_addrs; ni->v6_addrs = a6; }
a6->next = v6_head; v6_head = a6; }
}
ts = ts->next; }
}
/* Без изменений — освобождаем временные списки и выходим. */
if (sockmeta4_list_equal(ni->v4_sock_meta, m4_head) && addr4_list_equal(ni->v4_addrs, v4_head) &&
sockmeta6_list_equal(ni->v6_sock_meta, m6_head) && addr6_list_equal(ni->v6_addrs, v6_head)) {
free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, m4_head);
free_v4_addr_list(instance->topo_groups->v4_addr_pool, v4_head);
free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, m6_head);
free_v6_addr_list(instance->topo_groups->v6_addr_pool, v6_head);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my addresses unchanged (ver=%d), skip update", ni->ver);
return 0;
}
free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, ni->v4_sock_meta); ni->v4_sock_meta = m4_head;
free_v4_addr_list(instance->topo_groups->v4_addr_pool, ni->v4_addrs); ni->v4_addrs = v4_head;
free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, ni->v6_sock_meta); ni->v6_sock_meta = m6_head;
free_v6_addr_list(instance->topo_groups->v6_addr_pool, ni->v6_addrs); ni->v6_addrs = v6_head;
ni->ver = (ni->ver % 255) + 1;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "my addresses updated, new ver=%d", ni->ver);
topo_node_sign_self(instance, ni);
if (instance->topo_sqlite_db && topo_node_sqlite_member_exists(instance->topo_sqlite_db, instance->node_id))
if (instance->topo_sqlite_db) {
time_t now_sec = ntp_time_get_seconds(instance);
topo_node_sqlite_node_update_verified(instance->topo_sqlite_db, instance->node_id,
instance->name, instance->my_keys.public_key, instance->my_ed25519_pubkey,
(uint64_t)now_sec, now_sec);
topo_node_sqlite_addrs_put(instance->topo_sqlite_db, instance->node_id, ni);
}
struct ll_entry* ge = instance->topo_groups->group_list->head;
while (ge) {
@ -995,13 +1031,14 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
}
ge = ge->next;
}
return 1;
}
void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg) {
(void)arg;
if (!sock || !sock->instance) return;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "socket %s changed (event=0x%x), updating nodeinfo", sock->name, event);
topo_node_update_my_addresses(sock->instance);
if (sock->instance->topo_sqlite_db)
int changed = topo_node_update_my_addresses(sock->instance);
if (changed > 0 && sock->instance->topo_sqlite_db)
topo_node_sqlite_nodeinfo_updated(sock->instance->topo_sqlite_db, sock->instance->node_id);
}

2
src/routing_layer/topo_node.h

@ -198,7 +198,7 @@ int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t
uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count, uint16_t* out_rtt);
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);
void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance);
int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance); /* 1=changed, 0=unchanged, -1=error */
void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg);
void topo_node_dump_all(struct TOPO_GROUP* group);
int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size);

10
src/routing_layer/topo_node_sqlite.c

@ -111,7 +111,10 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE*
sqlite3_bind_int(is, 7, (int)cfg);
sqlite3_bind_int(is, 8, (int)nat);
sqlite3_bind_int(is, 9, (int)a->socket_id);
sqlite3_step(is); sqlite3_reset(is);
if (sqlite3_step(is) != SQLITE_DONE)
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=4 sock=%d: %s",
(unsigned long long)node_id, a->socket_id, sqlite3_errmsg(db));
sqlite3_reset(is);
written++;
}
for (const struct TOPO_ADDR6* a = ni->v6_addrs; a; a = a->next) {
@ -127,7 +130,10 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE*
sqlite3_bind_int(is, 7, (int)cfg);
sqlite3_bind_int(is, 8, (int)nat);
sqlite3_bind_int(is, 9, (int)a->socket_id);
sqlite3_step(is); sqlite3_reset(is);
if (sqlite3_step(is) != SQLITE_DONE)
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=6 sock=%d: %s",
(unsigned long long)node_id, a->socket_id, sqlite3_errmsg(db));
sqlite3_reset(is);
written++;
}
sqlite3_finalize(is);

Loading…
Cancel
Save