Browse Source

topo_group: migrate CONN_MGR entries to ll_queue, simplify TOPO_GROUP_NODE

- CONN_MGR entries: replace dynamic array with ll_queue (hash=256, key_size=8)
  cm_find_entry now O(1) via queue_find_data_by_index
  cm_ensure_entry uses queue_entry_new + queue_data_put_with_index
  cm_entry_cleanup removes from queue and frees entry

- TOPO_GROUP_NODE: replace handles** array with single CONN_MGR_HANDLE*
  remove topo_node_handle_register/unregister/count — dead indirection
  nq->handle = h / nq->handle = NULL in all call sites

- TOPO_GROUP_NODE: remove dead fields dirty, best_socket, tranzit_data/count
  dirty replaced with local variable in nat_detection.c
  tranzit_data removed from wire format (TOPOMSG_NODE), serialize, deserialize
  update struct layout: one field per line with comments
topo_upd
evgeny 2 months ago
parent
commit
c35135cb52
  1. 97
      src/routing_layer/conn_mgr_core.c
  2. 4
      src/routing_layer/conn_mgr_indirect.c
  3. 6
      src/routing_layer/conn_mgr_priv.h
  4. 7
      src/routing_layer/nat_detection.c
  5. 12
      src/routing_layer/topo_group.c
  6. 6
      src/routing_layer/topo_group.h
  7. 8
      src/routing_layer/topo_group_connect.c
  8. 43
      src/routing_layer/topo_node.c
  9. 34
      src/routing_layer/topo_node.h
  10. 7
      src/routing_layer/topo_recovery.c
  11. 2
      src/transport_layer/etcp_connections.c

97
src/routing_layer/conn_mgr_core.c

@ -154,24 +154,20 @@ int cm_has_local_addr(struct CONN_MGR* mgr, struct TOPO_GROUP_NODE* nq) {
/* ═══════ entries ═══════ */
struct CONN_MGR_ENTRY* cm_find_entry(struct CONN_MGR* mgr, uint64_t node_id) {
for (size_t i = 0; i < mgr->entry_count; i++)
if (mgr->entries[i].node_id == node_id && mgr->entries[i].state != CONN_MGR_STATE_DISCONNECTED)
return &mgr->entries[i];
return NULL;
if (!mgr || !mgr->entries) return NULL;
struct ll_entry* e = queue_find_data_by_index(mgr->entries, &node_id);
return e ? (struct CONN_MGR_ENTRY*)e : NULL;
}
struct CONN_MGR_ENTRY* cm_ensure_entry(struct CONN_MGR* mgr, uint64_t node_id) {
struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, node_id);
if (e) return e;
if (mgr->entry_count >= mgr->entry_capacity) {
size_t cap = mgr->entry_capacity * 2;
struct CONN_MGR_ENTRY* ne = u_realloc(mgr->entries, cap * sizeof(struct CONN_MGR_ENTRY));
if (!ne) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: realloc entries failed"); return NULL; }
memset(ne + mgr->entry_capacity, 0, (cap - mgr->entry_capacity) * sizeof(struct CONN_MGR_ENTRY));
mgr->entries = ne; mgr->entry_capacity = cap;
}
struct CONN_MGR_ENTRY* new_e = &mgr->entries[mgr->entry_count++];
memset(new_e, 0, sizeof(*new_e)); new_e->node_id = node_id; new_e->mgr = mgr;
struct ll_entry* qe = queue_entry_new(sizeof(struct CONN_MGR_ENTRY) - sizeof(struct ll_entry));
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: queue_entry_new failed"); return NULL; }
struct CONN_MGR_ENTRY* new_e = (struct CONN_MGR_ENTRY*)qe;
memset((uint8_t*)new_e + sizeof(struct ll_entry), 0, sizeof(*new_e) - sizeof(struct ll_entry));
new_e->node_id = node_id; new_e->mgr = mgr;
queue_data_put_with_index(mgr->entries, qe);
return new_e;
}
@ -192,6 +188,9 @@ void cm_entry_cleanup(struct CONN_MGR_ENTRY* entry) {
entry->handles = NULL;
cm_clear_nodeinfo(entry->mgr, entry->node_id);
entry->state = CONN_MGR_STATE_DISCONNECTED; entry->conn_type = CONN_TYPE_NONE;
{ struct CONN_MGR* m = entry->mgr;
queue_remove_data(m->entries, &entry->ll);
queue_entry_free(&entry->ll); }
}
void cm_clear_nodeinfo(struct CONN_MGR* mgr, uint64_t node_id) {
@ -276,7 +275,7 @@ void cm_invite_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_e
switch (ncd_ev) {
case NCD_EVENT_UP:
if (inv->state != CM_INVITE_CONNECTING) return;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite INIT OK node=0x%016llx", (unsigned long long)inv->node_id);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite INIT OK node=0x%016llx group=0x%016llx", (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id);
inv->state = CM_INVITE_WAIT_INFO;
{ struct CM_INVITE_REQ req; memset(&req, 0, sizeof(req));
req.cmd = ETCP_RT_ID_CONN_MGR; req.subcmd = CM_SUBCMD_INVITE_INFO_REQ; req.group_id = inv->mgr->group->group_id;
@ -306,8 +305,10 @@ struct CONN_MGR* conn_mgr_init(struct TOPO_GROUP* group) {
struct CONN_MGR* mgr = u_calloc(1, sizeof(struct CONN_MGR));
if (!mgr) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "alloc failed"); return NULL; }
mgr->instance = group->instance; mgr->group = group;
mgr->entry_capacity = 8; mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS;
mgr->entries = u_calloc(mgr->entry_capacity, sizeof(struct CONN_MGR_ENTRY));
mgr->direct_timeout_ms = CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS;
mgr->entries = queue_new(mgr->instance->ua, 256,
offsetof(struct CONN_MGR_ENTRY, node_id) - sizeof(struct ll_entry),
8, "conn_mgr_entries");
if (!mgr->entries) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "entries alloc failed"); u_free(mgr); return NULL; }
mgr->bg_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_BG_PING_INTERVAL_TB, mgr, cm_bg_ping_timer_cb, "conn_mgr_bg_ping");
mgr->candidate_ping_timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_CANDIDATE_PING_TB, mgr, cm_candidate_ping_timer_cb, "conn_mgr_cand_ping");
@ -326,16 +327,18 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) {
/* Отвязываемся от etcp_router (если ещё не отвязаны) */
etcp_router_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR);
etcp_unbind(mgr->instance, ETCP_RT_ID_CONN_MGR);
for (size_t i = 0; i < mgr->entry_count; i++) cm_entry_cleanup(&mgr->entries[i]);
while (mgr->invite_list) cm_invite_fail(mgr->invite_list);
while (mgr->exchange_pending) {
struct cm_exchange_pending* ep = mgr->exchange_pending; mgr->exchange_pending = ep->next;
if (ep->timer) uasync_cancel_timeout(mgr->instance->ua, ep->timer);
u_free(ep);
}
while (mgr->reverse_pending) { struct cm_reverse_pending* rp = mgr->reverse_pending; mgr->reverse_pending = rp->next; u_free(rp); }
u_free(mgr->entries); mgr->initialized = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: destroyed, entries=%zu", (size_t)mgr->entry_count);
{ size_t ec = queue_entry_count(mgr->entries);
struct ll_entry* e = mgr->entries->head;
while (e) { struct ll_entry* next = e->next; struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)e; cm_entry_cleanup(entry); e = next; }
while (mgr->invite_list) cm_invite_fail(mgr->invite_list);
while (mgr->exchange_pending) {
struct cm_exchange_pending* ep = mgr->exchange_pending; mgr->exchange_pending = ep->next;
if (ep->timer) uasync_cancel_timeout(mgr->instance->ua, ep->timer);
u_free(ep);
}
while (mgr->reverse_pending) { struct cm_reverse_pending* rp = mgr->reverse_pending; mgr->reverse_pending = rp->next; u_free(rp); }
queue_free(mgr->entries); mgr->initialized = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: destroyed, entries=%zu", ec); }
u_free(mgr);
}
@ -743,6 +746,9 @@ void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_
if (inv->mgr->instance->topo_sqlite_db) topo_node_sqlite_node_put(inv->mgr->instance->topo_sqlite_db, inv->mgr->instance->topo_groups, inv->temp_nq, (time_t)(get_time_tb()/10000));
}
if (inv->overall_timer) { uasync_cancel_timeout(inv->mgr->instance->ua, inv->overall_timer); inv->overall_timer = NULL; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite JOIN OK node=0x%016llx group=0x%016llx name=%.*s",
(unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id,
resp->node_name_len, resp->node_name_len ? (const char*)resp->node_name : "");
if (inv->cb) inv->cb(inv->handle, inv->node_id, inv->mgr->group->group_id, CONN_EVENT_JOIN, inv->cb_arg);
cm_invite_cleanup(inv);
}
@ -764,7 +770,17 @@ void cm_invite_fail(struct cm_invite_pending* inv) {
if (e) queue_remove_data(inv->mgr->instance->tcp_connections, e);
stcp_link_close(inv->tcp_link); inv->tcp_link = NULL;
}
if (inv->ncd_handle) { node_conn_direct_close(inv->ncd_handle); inv->ncd_handle = NULL; }
if (inv->ncd_handle) {
/* проверяем: есть ли зарегистрированные хендлы на этом соединении */
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(inv->mgr->group, inv->node_id);
if (nq && nq->handle != NULL) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite fail node=0x%016llx — handle registered, skip NCD close",
(unsigned long long)inv->node_id);
} else {
node_conn_direct_close(inv->ncd_handle);
}
inv->ncd_handle = NULL;
}
if (inv->temp_nq) {
queue_remove_data(inv->mgr->group->nodes, &inv->temp_nq->ll);
topo_nodeq_free_group_fields(inv->mgr->instance->topo_groups, inv->temp_nq);
@ -788,7 +804,7 @@ static void cm_tcp_ready_cb(struct stcp_link* link, void* arg) {
if (!qe) { cm_invite_fail(inv); return; }
{ struct tcp_conn_entry* te = (struct tcp_conn_entry*)qe->data; te->node_id=inv->node_id; te->link=link; te->etcp_conn=etcp; }
queue_data_put_with_index(inv->mgr->instance->tcp_connections, qe);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP invite ready node=0x%016llx", (unsigned long long)inv->node_id);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP invite ready node=0x%016llx group=0x%016llx", (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id);
inv->state = CM_INVITE_WAIT_INFO;
{ struct CM_INVITE_REQ req; memset(&req,0,sizeof(req));
req.cmd=ETCP_RT_ID_CONN_MGR; req.subcmd=CM_SUBCMD_INVITE_INFO_REQ; req.group_id=inv->mgr->group->group_id;
@ -816,7 +832,11 @@ static int cm_invite_tcp_connect(struct cm_invite_pending* inv, const struct TOP
struct stcp_link_config cfg = {.ua=inv->mgr->instance->ua, .my_keys=&inv->mgr->instance->my_keys, .inst=inv->mgr->instance,
.peer_pubkey=ni->public_key, .peer_pubkey_mode=0, .remote_addr=&sa, .remote_port=a->port};
inv->tcp_link = stcp_link_connect(&cfg);
if (inv->tcp_link) { stcp_link_set_on_ready(inv->tcp_link, cm_tcp_ready_cb, inv); stcp_link_set_on_close(inv->tcp_link, cm_tcp_close_cb, inv); return 1; }
if (inv->tcp_link) { stcp_link_set_on_ready(inv->tcp_link, cm_tcp_ready_cb, inv); stcp_link_set_on_close(inv->tcp_link, cm_tcp_close_cb, inv);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP connect %d.%d.%d.%d:%d node=0x%016llx group=0x%016llx",
a->addr[0], a->addr[1], a->addr[2], a->addr[3], (int)a->port,
(unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id);
return 1; }
return 0;
}
for (const struct TOPO_ADDR6* a6 = ni->v6_addrs; a6; a6 = a6->next) {
@ -831,7 +851,10 @@ static int cm_invite_tcp_connect(struct cm_invite_pending* inv, const struct TOP
struct stcp_link_config cfg = {.ua=inv->mgr->instance->ua, .my_keys=&inv->mgr->instance->my_keys, .inst=inv->mgr->instance,
.peer_pubkey=ni->public_key, .peer_pubkey_mode=0, .remote_addr=&sa, .remote_port=a6->port};
inv->tcp_link = stcp_link_connect(&cfg);
if (inv->tcp_link) { stcp_link_set_on_ready(inv->tcp_link, cm_tcp_ready_cb, inv); stcp_link_set_on_close(inv->tcp_link, cm_tcp_close_cb, inv); return 1; }
if (inv->tcp_link) { stcp_link_set_on_ready(inv->tcp_link, cm_tcp_ready_cb, inv); stcp_link_set_on_close(inv->tcp_link, cm_tcp_close_cb, inv);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP6 connect port=%d node=0x%016llx group=0x%016llx",
(int)a6->port, (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id);
return 1; }
return 0;
}
return 0;
@ -905,7 +928,17 @@ int conn_mgr_open_invite(struct UTUN_INSTANCE* inst,
if (!udp && !tcp) { cm_invite_fail(inv); if (out_handle) *out_handle=NULL; return -1; }
inv->overall_timer = uasync_set_timeout(inst->ua, CM_INVITE_DEFAULT_TIMEOUT_MS*10, inv, cm_invite_overall_timeout, "conn_mgr_invite");
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite started node=0x%016llx group=0x%016llx",
(unsigned long long)nid, (unsigned long long)group_id);
{
char addrs[256] = ""; int aoff = 0;
for (const struct TOPO_ADDR4* a = gni->v4_addrs; a; a = a->next)
aoff += snprintf(addrs + aoff, sizeof(addrs) - (size_t)aoff, " %d.%d.%d.%d:%d(%s)",
a->addr[0], a->addr[1], a->addr[2], a->addr[3], (int)a->port,
a->protocol == TOPO_PROTO_UDP ? "UDP" : a->protocol == TOPO_PROTO_TCP ? "TCP" : "?");
for (const struct TOPO_ADDR6* a6 = gni->v6_addrs; a6; a6 = a6->next)
aoff += snprintf(addrs + aoff, sizeof(addrs) - (size_t)aoff, " [v6]:%d(%s)", (int)a6->port,
a6->protocol == TOPO_PROTO_UDP ? "UDP" : a6->protocol == TOPO_PROTO_TCP ? "TCP" : "?");
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite started node=0x%016llx group=0x%016llx udp=%d tcp=%d%s",
(unsigned long long)nid, (unsigned long long)group_id, udp, tcp, addrs);
}
return 0;
}

4
src/routing_layer/conn_mgr_indirect.c

@ -109,8 +109,8 @@ void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, s
void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
const struct CM_INTERM_SEL* sel=(const struct CM_INTERM_SEL*)data;
if(len<offsetof(struct CM_INTERM_SEL,selected))return;
struct CONN_MGR_ENTRY* entry=NULL;
for(size_t i=0;i<mgr->entry_count;i++) if(mgr->entries[i].state==CONN_MGR_STATE_CONNECTING){entry=&mgr->entries[i];break;}
struct CONN_MGR_ENTRY* entry = NULL; struct ll_entry* ee = mgr->entries->head;
while (ee) { struct CONN_MGR_ENTRY* e = (struct CONN_MGR_ENTRY*)ee; if (e->state == CONN_MGR_STATE_CONNECTING) { entry = e; break; } ee = ee->next; }
if(!entry)return;
uint8_t c=sel->count>CONN_MGR_MAX_INTERMEDIARIES?CONN_MGR_MAX_INTERMEDIARIES:sel->count;
for(uint8_t i=0;i<c;i++)entry->intermediaries[i]=sel->selected[i].node_id;

6
src/routing_layer/conn_mgr_priv.h

@ -121,7 +121,9 @@ struct CONN_MGR_HANDLE {
};
struct CONN_MGR_ENTRY {
uint64_t node_id; uint8_t state, conn_type;
struct ll_entry ll; /* ll_queue entry — первый элемент */
uint64_t node_id; /* хеш-ключ 8B, offset от data = 0 */
uint8_t state, conn_type;
uint32_t idle_timeout_ms; uint64_t last_traffic_tb; void* idle_timer;
struct CONN_MGR_HANDLE* handles; struct NODE_CONN_DIRECT* ncd_handle;
uint64_t intermediaries[3]; uint8_t intermediariy_count, rr_idx;
@ -132,7 +134,7 @@ struct CONN_MGR_ENTRY {
struct CONN_MGR {
struct UTUN_INSTANCE* instance; struct TOPO_GROUP* group;
struct CONN_MGR_ENTRY* entries; size_t entry_count, entry_capacity;
struct ll_queue* entries;
void *bg_ping_timer, *candidate_ping_timer;
size_t bg_ping_cursor; uint64_t bg_ping_cycle_start_tb;
struct CONN_MGR_CANDIDATE best_candidates[CONN_MGR_MAX_CANDIDATES]; uint8_t best_candidate_count;

7
src/routing_layer/nat_detection.c

@ -167,9 +167,10 @@ static void nat_detection_handle_nat_info(struct NAT_DETECTION* nd, struct TOPO_
if (data_changed) {
int prev_v4a = topo_list_count((struct _topo_head*)ni->v4_addrs);
int dirty = 0;
topo_group_update_my_nodeinfo(group->instance, group);
if (topo_list_count((struct _topo_head*)ni->v4_addrs) != prev_v4a) {
group->local_node->dirty = 1;
dirty = 1;
ni->ver = (ni->ver % 255) + 1;
group->local_node->last_ver = ni->ver;
}
@ -180,7 +181,7 @@ static void nat_detection_handle_nat_info(struct NAT_DETECTION* nd, struct TOPO_
uint32_t a_ip; memcpy(&a_ip, a->addr, 4);
if (a_ip == nat_ip && a->port == nat_port) {
a->type = TOPO_ADDR_NAT; a->socket_id = socket_id | 1;
group->local_node->dirty = 1;
dirty = 1;
ni->ver = (ni->ver % 255) + 1;
group->local_node->last_ver = ni->ver;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "NAT_INFO matched interface addr, updated entry type to NAT sock=%d", socket_id);
@ -190,7 +191,7 @@ static void nat_detection_handle_nat_info(struct NAT_DETECTION* nd, struct TOPO_
a = a->next;
}
}
if (group->local_node->dirty && group->senders_list) {
if (dirty && group->senders_list) {
struct ll_entry* se = group->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;

12
src/routing_layer/topo_group.c

@ -439,7 +439,6 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
if (key != conn->peer_node_id) { topo_recovery_add_node(group, nq, conn->peer_node_id); cascaded++; }
if (group->group_type != TOPO_GROUP_TYPE_CHAT && rt) route_delete(rt, nq);
if (conn->instance && conn->instance->control_srv) control_server_notify_node_removed(conn->instance->control_srv, key);
nq->dirty = 1;
route_connectivity_cancel_node(conn->instance, nq);
if (nq->paths) { queue_free(nq->paths); nq->paths = NULL; }
topo_nodeq_free_group_fields(group->instance->topo_groups, nq);
@ -600,10 +599,9 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
size_t ser_len = len - 2;
struct TOPO_NODE* new_ni = NULL;
struct TOPO_NODESUBNETS* new_subnets = NULL;
struct TOPOMSG_TRANZIT* new_tranzit = NULL; uint8_t new_tranzit_count = 0;
uint64_t* new_hop_list = NULL; uint8_t new_hop_count = 0;
uint16_t incoming_cumulative_rtt = 0;
if (topo_node_deserialize(group, ser_data, ser_len, &new_ni, &new_subnets, &new_tranzit, &new_tranzit_count, &new_hop_list, &new_hop_count, &incoming_cumulative_rtt) != 0) {
if (topo_node_deserialize(group, ser_data, ser_len, &new_ni, &new_subnets, &new_hop_list, &new_hop_count, &incoming_cumulative_rtt) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "NODEINFO deserialize failed from %s nid=%016llx", from->log_name, (unsigned long long)node_id);
if (nodeinfo1) { queue_remove_data(group->nodes, &nodeinfo1->ll); queue_free(paths); queue_entry_free(&nodeinfo1->ll); }
return -1;
@ -621,7 +619,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "NODEINFO x25519_self_sig VERIFY FAIL node=%016llx ed_pubkey=%016llx... from=%s — rejecting as forgery",
(unsigned long long)node_id, ekchk, from->log_name);
topo_node_destroy(group->instance->topo_groups, new_ni);
u_free(new_subnets); u_free(new_tranzit); u_free(new_hop_list);
u_free(new_subnets); u_free(new_hop_list);
if (nodeinfo1) { queue_remove_data(group->nodes, &nodeinfo1->ll); queue_free(paths); queue_entry_free(&nodeinfo1->ll); }
return -1;
} else {
@ -640,13 +638,11 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
nodeinfo1->node_id = new_ni ? new_ni->node_id : 0;
nodeinfo1->subnets = new_subnets;
nodeinfo1->cumulative_rtt = incoming_cumulative_rtt;
u_free(nodeinfo1->tranzit_data); nodeinfo1->tranzit_data = new_tranzit;
nodeinfo1->tranzit_count = new_tranzit_count;
u_free(nodeinfo1->hop_list); nodeinfo1->hop_list = new_hop_list;
nodeinfo1->hop_count = new_hop_count;
} else {
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_GROUP_NODE));
if (!qe) { topo_node_destroy(group->instance->topo_groups, new_ni); u_free(new_tranzit); u_free(new_hop_list); u_free(new_subnets); return -1; }
if (!qe) { topo_node_destroy(group->instance->topo_groups, new_ni); u_free(new_hop_list); u_free(new_subnets); return -1; }
nodeinfo1 = (struct TOPO_GROUP_NODE*)qe;
memset((uint8_t*)nodeinfo1 + sizeof(struct ll_entry), 0, sizeof(*nodeinfo1) - sizeof(struct ll_entry));
{ struct TOPO_NODE* stored = topo_node_registry_store(group->instance->topo_groups, new_ni);
@ -658,7 +654,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
nodeinfo1->node_id = new_ni ? new_ni->node_id : 0;
nodeinfo1->subnets = new_subnets;
nodeinfo1->cumulative_rtt = incoming_cumulative_rtt;
nodeinfo1->tranzit_data = new_tranzit; nodeinfo1->tranzit_count = new_tranzit_count;
nodeinfo1->hop_list = new_hop_list; nodeinfo1->hop_count = new_hop_count;
nodeinfo1->connectivity.probe_status = PROBE_STATUS_NONE;
nodeinfo1->connectivity.interface_status = PROBE_RESULT_UNKNOWN;
@ -747,7 +742,6 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
if (ret == 1) {
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance && group->instance->rt) route_delete(group->instance->rt, nq);
if (group->instance->control_srv) control_server_notify_node_removed(group->instance->control_srv, node_id);
nq->dirty = 1;
route_connectivity_cancel_node(group->instance, nq);
if (nq->paths) { queue_free(nq->paths); nq->paths = NULL; }
topo_nodeq_free_group_fields(group->instance->topo_groups, nq);

6
src/routing_layer/topo_group.h

@ -146,9 +146,9 @@ struct TOPO_GROUP {
uint64_t group_id; // уникальный идентификатор группы
uint8_t group_type; // TOPO_GROUP_TYPE_*
struct UTUN_INSTANCE* instance;
struct ll_queue* senders_list;
struct ll_queue* nodes;
struct TOPO_GROUP_NODE* local_node;
struct ll_queue* senders_list; // TOPO_GROUP_CONN_ITEM{ll_entry,ETCP_CONN*} — активные BGP-пиры, без хеш-индекса
struct ll_queue* nodes; // TOPO_GROUP_NODE{ll_entry,node_id,paths,...} — per-group узлы, хеш-индекс по node_id(8B)
struct TOPO_GROUP_NODE* local_node; // свой узел в этой группе
uint8_t ed25519_public_key[SC_PUBKEY_SIZE];
char channel_id[64]; // channel_id для групп типа CHAT
struct CONN_MGR* conn_mgr; // менеджер соединений для этой группы

8
src/routing_layer/topo_group_connect.c

@ -201,6 +201,14 @@ static void tgc_callback(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t g
int ok = (event == CONN_EVENT_UP);
if (ok) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(gc->group, node_id);
if (nq) nq->handle = h;
} else if (event == CONN_EVENT_TIMEOUT) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(gc->group, node_id);
if (nq) nq->handle = NULL;
}
switch (gc->phase) {
case TGC_PHASE_ONE:
gc->pending--;

43
src/routing_layer/topo_node.c

@ -87,8 +87,8 @@ void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_
u_free(nq->subnets);
nq->subnets = NULL;
}
u_free(nq->tranzit_data); nq->tranzit_data = NULL; nq->tranzit_count = 0;
u_free(nq->hop_list); nq->hop_list = NULL; nq->hop_count = 0;
nq->handle = NULL;
}
// ===== глобальный реестр TOPO_NODE =====
@ -145,9 +145,8 @@ int topo_node_dyn_size(const struct TOPOMSG_NODE* msg) {
+ msg->local_v6_sockets * sizeof(struct TOPOMSG_SOCKMETA6)
+ msg->local_v6_addrs * sizeof(struct TOPOMSG_ADDR6)
+ ((msg->flags & TOPO_FLAG_SEND_SUBNETS)
? (msg->local_v4_subnets * sizeof(struct TOPOMSG_SUBNET4) + msg->local_v6_subnets * sizeof(struct TOPOMSG_SUBNET6))
: 0)
+ msg->tranzit_nodes * sizeof(struct TOPOMSG_TRANZIT)
? (msg->local_v4_subnets * sizeof(struct TOPOMSG_SUBNET4) + msg->local_v6_subnets * sizeof(struct TOPOMSG_SUBNET6))
: 0)
+ msg->hop_count * 8;
}
@ -170,7 +169,6 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_
msg.local_v6_addrs = topo_list_count((struct _topo_head*)ni->v6_addrs);
msg.local_v4_subnets = nq->subnets ? topo_list_count((struct _topo_head*)nq->subnets->v4_subnets) : 0;
msg.local_v6_subnets = nq->subnets ? topo_list_count((struct _topo_head*)nq->subnets->v6_subnets) : 0;
msg.tranzit_nodes = nq->tranzit_count;
msg.hop_count = nq->hop_count;
msg.cumulative_rtt = cumulative_rtt;
@ -210,10 +208,6 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_
}}
}
}
if (msg.tranzit_nodes && nq->tranzit_data) {
size_t sz = msg.tranzit_nodes * sizeof(struct TOPOMSG_TRANZIT);
memcpy(dp, nq->tranzit_data, sz); dp += sz;
}
if (msg.hop_count && nq->hop_list) {
size_t sz = msg.hop_count * 8;
memcpy(dp, nq->hop_list, sz); dp += sz;
@ -224,7 +218,6 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_
// ===== deserialization =====
int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t len, struct TOPO_NODE** out_ni, struct TOPO_NODESUBNETS** out_subnets,
struct TOPOMSG_TRANZIT** out_tranzit, uint8_t* out_tranzit_count,
uint64_t** out_hop_list, uint8_t* out_hop_count,
uint16_t* out_cumulative_rtt) {
if (!group || !data || len < TOPOMSG_NODE_HDR_SIZE || !out_ni) return -1;
@ -307,13 +300,6 @@ int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t
if (out_subnets) *out_subnets = NULL;
}
if (out_tranzit && msg->tranzit_nodes > 0) {
size_t tsz = msg->tranzit_nodes * sizeof(struct TOPOMSG_TRANZIT);
*out_tranzit = u_malloc(tsz);
if (!*out_tranzit) { topo_node_free_raw(group->instance->topo_groups, ni); return -1; }
memcpy(*out_tranzit, dp, tsz); dp += tsz;
*out_tranzit_count = msg->tranzit_nodes;
}
if (out_hop_list && msg->hop_count > 0) {
size_t hsz = msg->hop_count * 8;
*out_hop_list = u_malloc(hsz);
@ -419,12 +405,10 @@ void topo_node_dump_all(struct TOPO_GROUP* group) {
int is_local = (group->local_node == nq);
const char* name_s = ni->node_name ? ni->node_name : "";
DEBUG_INFO(DEBUG_CATEGORY_BGP, "--- Node #%d: id=%016llx name=\"%s\" ver=%u dirty=%u grp=%llu%s ---",
idx, (unsigned long long)ni->node_id, name_s, ni->ver, nq->dirty,
DEBUG_INFO(DEBUG_CATEGORY_BGP, "--- Node #%d: id=%016llx name=\"%s\" ver=%u grp=%llu%s ---",
idx, (unsigned long long)ni->node_id, name_s, ni->ver,
(unsigned long long)ni->group_id, is_local ? " [LOCAL]" : "");
if (nq->best_socket) DEBUG_INFO(DEBUG_CATEGORY_BGP, " best_socket=%p", (void*)nq->best_socket);
log_dump(DEBUG_LEVEL_INFO, DEBUG_CATEGORY_BGP, " pubkey ", ni->public_key, SC_PUBKEY_SIZE);
log_dump(DEBUG_LEVEL_INFO, DEBUG_CATEGORY_BGP, " ed25519", ni->ed25519_public_key, SC_PUBKEY_SIZE);
@ -460,11 +444,6 @@ void topo_node_dump_all(struct TOPO_GROUP* group) {
sub6 = sub6->next;
}}
}
for (int i = 0; i < nq->tranzit_count; i++) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, " tranzit[%d]: node=%016llx rtt=%u linkq=%u",
i, (unsigned long long)nq->tranzit_data[i].node_id, nq->tranzit_data[i].rtt, nq->tranzit_data[i].link_q);
}
int path_idx = 0;
if (nq->paths) {
struct ll_entry* pe = nq->paths->head;
@ -517,12 +496,10 @@ int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size) {
#define FMT_ADD(fmt, ...) pos += snprintf(buf + pos, buf_size > (size_t)pos ? buf_size - pos : 0, fmt, ##__VA_ARGS__)
FMT_ADD("--- Node #%d: id=%016llx name=\"%s\" ver=%u dirty=%u grp=%llu%s ---\n",
idx, (unsigned long long)ni->node_id, name_s, ni->ver, nq->dirty,
FMT_ADD("--- Node #%d: id=%016llx name=\"%s\" ver=%u grp=%llu%s ---\n",
idx, (unsigned long long)ni->node_id, name_s, ni->ver,
(unsigned long long)ni->group_id, is_local ? " [LOCAL]" : "");
if (nq->best_socket) FMT_ADD(" best_socket=%p\n", (void*)nq->best_socket);
FMT_ADD(" pubkey: "); for (int k = 0; k < SC_PUBKEY_SIZE; k++) FMT_ADD("%02x", ni->public_key[k]); FMT_ADD("\n");
FMT_ADD(" ed25519: "); for (int k = 0; k < SC_PUBKEY_SIZE; k++) FMT_ADD("%02x", ni->ed25519_public_key[k]); FMT_ADD("\n");
@ -552,10 +529,6 @@ int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size) {
FMT_ADD(" v6sub[%d]: %s/%d\n", r6i++, ip_to_str(sub6->addr, AF_INET6).str, sub6->prefix_length); sub6 = sub6->next;
}}
}
for (int i = 0; i < nq->tranzit_count; i++)
FMT_ADD(" tranzit[%d]: node=%016llx rtt=%u linkq=%u\n",
i, (unsigned long long)nq->tranzit_data[i].node_id, nq->tranzit_data[i].rtt, nq->tranzit_data[i].link_q);
if (nq->paths) {
struct ll_entry* pe = nq->paths->head; int pi = 0;
while (pe) {
@ -686,7 +659,6 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR
topo_node_registry_store(instance->topo_groups, ni);
lq->node_id = ni->node_id;
lq->hop_count = 0;
lq->dirty = 1;
lq->last_ver = ni->ver;
e_sock = instance->etcp_sockets;
@ -838,7 +810,6 @@ void topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) {
while (ge) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge;
if (g->local_node) {
g->local_node->dirty = 1;
if (g->senders_list) {
struct ll_entry* se = g->senders_list->head;
while (se) {

34
src/routing_layer/topo_node.h

@ -33,6 +33,7 @@ extern "C" {
struct TOPO_GROUP;
struct TOPO_GROUPS;
struct ETCP_SOCKET;
struct CONN_MGR_HANDLE;
struct ETCP_CONN;
struct UTUN_INSTANCE;
@ -97,14 +98,13 @@ struct TOPOMSG_SOCKMETA6 { uint8_t id, config_type, nat_type; } __attribute__((p
struct TOPOMSG_ADDR6 { uint8_t addr[16]; uint16_t port; uint8_t type, socket_id, protocol; } __attribute__((packed));
struct TOPOMSG_SUBNET4 { uint8_t addr[4]; uint8_t prefix_length; } __attribute__((packed));
struct TOPOMSG_SUBNET6 { uint8_t addr[16]; uint8_t prefix_length; } __attribute__((packed));
struct TOPOMSG_TRANZIT { uint64_t node_id; uint16_t rtt, link_q; } __attribute__((packed));
struct TOPOMSG_NODE {
uint8_t flags; uint64_t group_id, node_id; uint8_t ver;
uint8_t public_key[SC_PUBKEY_SIZE], ed25519_public_key[SC_PUBKEY_SIZE];
uint8_t x25519_self_sig[64];
uint8_t node_name_len, local_v4_sockets, local_v4_addrs, local_v6_sockets, local_v6_addrs;
uint8_t local_v4_subnets, local_v6_subnets, tranzit_nodes, hop_count;
uint8_t local_v4_subnets, local_v6_subnets, hop_count;
uint16_t cumulative_rtt;
} __attribute__((packed));
@ -146,21 +146,21 @@ struct TOPO_NODEPATH {
uint16_t cumulative_rtt;
};
/** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes. */
/** Узел в контексте группы (per-group данные). Хранится в TOPO_GROUP->nodes (ll_queue). */
struct TOPO_GROUP_NODE {
struct ll_entry ll;
uint64_t node_id;
struct TOPO_NODESUBNETS* subnets;
struct TOPOMSG_TRANZIT* tranzit_data;
uint64_t* hop_list;
uint8_t tranzit_count, hop_count;
uint16_t cumulative_rtt;
struct ll_queue* paths;
uint8_t dirty, last_ver, conn_mgr_type;
uint64_t conn_mgr_intermediaries[CONN_MGR_MAX_INTERMEDIARIES];
uint8_t conn_mgr_intermediariy_count;
struct TOPO_CONNECTIVITY connectivity;
struct ETCP_SOCKET* best_socket;
struct ll_entry ll; // queue entry, хеш-ключ 8 байт → node_id
uint64_t node_id; // идентификатор узла в этой группе
struct TOPO_NODESUBNETS* subnets; // подсети узла (v4/v6), для route_insert()
uint64_t* hop_list; // путь от нас к узлу (массив node_id)
uint8_t hop_count; // длина hop_list
uint16_t cumulative_rtt; // накопленный RTT от нас до узла
struct ll_queue* paths; // TOPO_NODEPATH{conn,hop_count,rtt} — маршруты
uint8_t last_ver; // последняя версия NODEINFO (защита от stale, BGP)
uint8_t conn_mgr_type; // CONN_TYPE_DIRECT/REVERSE/INDIRECT/NONE
uint64_t conn_mgr_intermediaries[CONN_MGR_MAX_INTERMEDIARIES]; // посредники для INDIRECT
uint8_t conn_mgr_intermediariy_count; // количество посредников
struct TOPO_CONNECTIVITY connectivity; // результаты ping-проб (interface/nat/real)
struct CONN_MGR_HANDLE* handle; // активный хендл соединения к узлу
};
// API — глобальный реестр TOPO_NODE
@ -177,7 +177,7 @@ void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_
int topo_node_dyn_size(const struct TOPOMSG_NODE* msg);
int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_t* out, size_t out_max, uint16_t cumulative_rtt);
int topo_node_deserialize(struct TOPO_GROUP* group, const uint8_t* data, size_t len, struct TOPO_NODE** out_ni, struct TOPO_NODESUBNETS** out_subnets,
struct TOPOMSG_TRANZIT** out_tranzit, uint8_t* out_tranzit_count, uint64_t** out_hop_list, uint8_t* out_hop_count,
uint64_t** out_hop_list, uint8_t* out_hop_count,
uint16_t* out_cumulative_rtt);
int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group);

7
src/routing_layer/topo_recovery.c

@ -106,10 +106,13 @@ static void topo_recovery_callback(struct CONN_MGR_HANDLE* h, uint64_t node_id,
int list = ctx->current_list;
for (size_t i = 0; i < ctx->count[list]; i++)
if (ctx->nodes[list][i].node_id == node_id) { ctx->nodes[list][i] = ctx->nodes[list][ctx->count[list] - 1]; ctx->count[list]--; break; }
if (event == CONN_EVENT_UP)
if (event == CONN_EVENT_UP) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(ctx->group, node_id);
if (nq) nq->handle = h;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx reconnected to %016llx", (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id);
else
} else {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "recovery: next=%016llx connect to %016llx failed: %d", (unsigned long long)ctx->next_hop_id, (unsigned long long)node_id, event);
}
ctx->current_node_id = 0;
topo_recovery_try_next(ctx);
}

2
src/transport_layer/etcp_connections.c

@ -1408,8 +1408,6 @@ static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pk
if (g->local_node && g->senders_list) {
struct TOPO_NODE* vni = topo_node_registry_find(g->instance->topo_groups, g->local_node->node_id);
if (vni) vni->ver = (vni->ver % 255) + 1;
g->local_node->dirty = 1;
g->local_node->dirty = 1;
struct ll_entry* se = g->senders_list->head;
while (se) {
struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data;

Loading…
Cancel
Save