|
|
|
|
@ -48,6 +48,20 @@ static void topo_group_send_table_request(struct TOPO_GROUP* group, struct ETCP_
|
|
|
|
|
if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void topo_group_send_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { |
|
|
|
|
if (!group || !conn) return; |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Sending join group request to %s grp=%016llx", conn->log_name, (unsigned long long)group->group_id); |
|
|
|
|
struct TOPOMSG_JOIN_GROUP* req = u_calloc(1, sizeof(struct TOPOMSG_JOIN_GROUP)); |
|
|
|
|
if (!req) return; |
|
|
|
|
req->cmd = ETCP_ID_TOPO_ENTRY; |
|
|
|
|
req->subcmd = TOPO_SUBCMD_JOIN_GROUP; |
|
|
|
|
req->group_id = group->group_id; |
|
|
|
|
struct ll_entry* e = queue_entry_new(0); |
|
|
|
|
if (!e) { u_free(req); return; } |
|
|
|
|
e->dgram = (uint8_t*)req; e->len = sizeof(struct TOPOMSG_JOIN_GROUP); |
|
|
|
|
if (etcp_send(conn, e) != 0) { u_free(req); queue_entry_free(e); } |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void topo_group_send_table_complete(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { |
|
|
|
|
if (!group || !conn) return; |
|
|
|
|
struct TOPOMSG_TABLE_REQ* req = u_calloc(1, sizeof(struct TOPOMSG_TABLE_REQ)); |
|
|
|
|
@ -65,6 +79,7 @@ static void topo_group_add_to_senders(struct TOPO_GROUP* group, struct ETCP_CONN
|
|
|
|
|
static bool topo_group_should_send_to(const struct TOPO_GROUP_NODE* nq, uint64_t target_id); |
|
|
|
|
static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn); |
|
|
|
|
static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETCP_CONN* conn); |
|
|
|
|
static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn); |
|
|
|
|
|
|
|
|
|
static void nodeinfo_dump_log(const uint8_t* data, size_t len) { |
|
|
|
|
if (!data || len < sizeof(struct TOPOMSG_NODEINFO_PKT)) return; |
|
|
|
|
@ -114,6 +129,7 @@ static const char* group_subcmd_name(uint8_t subcmd) {
|
|
|
|
|
case TOPO_SUBCMD_WITHDRAW: return "WITHDRAW"; |
|
|
|
|
case TOPO_SUBCMD_TABLE_COMPLETE: return "TABLE_COMPLETE"; |
|
|
|
|
case TOPO_SUBCMD_ERR_GROUP_MISMATCH: return "ERR_GROUP_MISMATCH"; |
|
|
|
|
case TOPO_SUBCMD_JOIN_GROUP: return "JOIN_GROUP"; |
|
|
|
|
default: return "?"; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
@ -164,6 +180,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)) { |
|
|
|
|
pkt_group_id = ((struct TOPOMSG_JOIN_GROUP*)data)->group_id; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
struct TOPO_GROUP* group = topo_groups_find(instance->topo_groups, pkt_group_id); |
|
|
|
|
@ -177,6 +195,7 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
|
|
|
|
|
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_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)) { |
|
|
|
|
@ -194,14 +213,77 @@ static void topo_group_receive_cbk(struct ETCP_CONN* from_conn, struct ll_entry*
|
|
|
|
|
// Init / Destroy / New / Remove conn
|
|
|
|
|
// ============================================================================
|
|
|
|
|
|
|
|
|
|
/* Создать TOPO_GROUP_NODE (если нет) + прямой путь через conn (hop=1). Без BGP.
|
|
|
|
|
Используется из topo_group_new_conn (add + init BGP) и из приёма TABLE_REQ (входящий запросил BGP). */ |
|
|
|
|
static void topo_group_conn_add_path(struct TOPO_GROUP* group, struct ETCP_CONN* conn, uint64_t node_id) { |
|
|
|
|
if (!group || !conn || !node_id || node_id == group->instance->node_id) return; |
|
|
|
|
|
|
|
|
|
/* узел в реестре — иначе загрузить из БД (nodes + node_addresses), чтобы форвардинг мог сериализовать */ |
|
|
|
|
{ struct TOPO_GROUPS* groups = group->instance->topo_groups; |
|
|
|
|
struct TOPO_NODE* ni = topo_node_registry_find(groups, node_id); |
|
|
|
|
if (!ni) { |
|
|
|
|
sqlite3* db = group->instance->topo_sqlite_db; |
|
|
|
|
ni = db ? topo_node_sqlite_node_load(db, groups, node_id) : NULL; |
|
|
|
|
if (ni) ni = topo_node_registry_store(groups, ni); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id); |
|
|
|
|
if (nq && nq->paths) { |
|
|
|
|
struct ll_entry* pe = nq->paths->head; |
|
|
|
|
while (pe) { if (((struct TOPO_NODEPATH*)pe)->conn == conn) return; pe = pe->next; } |
|
|
|
|
} |
|
|
|
|
if (!nq) { |
|
|
|
|
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_GROUP_NODE)); |
|
|
|
|
if (!qe) return; |
|
|
|
|
nq = (struct TOPO_GROUP_NODE*)qe; |
|
|
|
|
memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry)); |
|
|
|
|
nq->node_id = node_id; |
|
|
|
|
queue_data_put_with_index(group->nodes, &nq->ll); |
|
|
|
|
} |
|
|
|
|
uint64_t hop[1] = { node_id }; |
|
|
|
|
topo_group_add_path(nq, conn, hop, 1, 0); |
|
|
|
|
nq->conn_presence |= NCONN_DIRECT; |
|
|
|
|
nq->conn_up |= NCONN_DIRECT; |
|
|
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "conn_add_path: node=0x%016llx grp=%016llx ch=%s conn=%s", |
|
|
|
|
(unsigned long long)node_id, (unsigned long long)group->group_id, group->channel_id, conn->log_name); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg) { |
|
|
|
|
struct TOPO_GROUPS* groups = (struct TOPO_GROUPS*)arg; |
|
|
|
|
if (!conn || !groups) return; |
|
|
|
|
struct ll_entry* fe = groups->group_list->head; |
|
|
|
|
while (fe) { struct TOPO_GROUP* g = (struct TOPO_GROUP*)fe; fe = fe->next; |
|
|
|
|
if (status == ETCP_CONN_STATUS_UP) topo_group_new_conn(g, conn); |
|
|
|
|
if (status == ETCP_CONN_STATUS_DOWN) topo_group_remove_conn(g, conn); |
|
|
|
|
if (status == ETCP_CONN_STATUS_DELETE) topo_group_remove_conn(g, conn); |
|
|
|
|
|
|
|
|
|
if (status == ETCP_CONN_STATUS_UP && conn->peer_node_id) { |
|
|
|
|
/* non-CHAT (UTUN): узлы из конфига (clients) + явные подключения (etcp_connect). */ |
|
|
|
|
int explicit_conn = conn->bgp_ready_cbk != NULL; |
|
|
|
|
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) { |
|
|
|
|
int is_client = conn->instance->config && config_peer_in_clients(conn->instance->config, conn->crypto_ctx.peer_public_key); |
|
|
|
|
if (is_client || explicit_conn) { |
|
|
|
|
topo_group_new_conn(g, conn); /* добавить + инициировать BGP */ |
|
|
|
|
topo_group_send_join_group(g, conn); /* запросить членство у пира */ |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
/* CHAT: только члены каналов (peers_* в БД). */ |
|
|
|
|
sqlite3* db = conn->instance ? conn->instance->topo_sqlite_db : NULL; |
|
|
|
|
uint64_t* chs = NULL; int chn = 0; |
|
|
|
|
if (db && topo_node_sqlite_get_member_channels(db, conn->peer_node_id, &chs, &chn) == 0 && chn > 0) { |
|
|
|
|
for (int i = 0; i < chn; i++) { |
|
|
|
|
struct TOPO_GROUP* cg = topo_groups_find(groups, chs[i]); |
|
|
|
|
if (cg && cg->group_type == TOPO_GROUP_TYPE_CHAT) |
|
|
|
|
topo_group_new_conn(cg, conn); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
u_free(chs); |
|
|
|
|
} 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); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -467,9 +549,21 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
|
|
|
|
|
if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance is NULL"); return; } |
|
|
|
|
if (!conn->instance->rt) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance->rt is NULL"); return; } |
|
|
|
|
|
|
|
|
|
topo_recovery_cancel_for_node(group, conn->peer_node_id); |
|
|
|
|
|
|
|
|
|
/* создать узел + прямой путь (сразу, не дожидаясь NODEINFO; дедуп внутри по conn).
|
|
|
|
|
Делаем ДО проверки senders_list: conn может уже быть в senders_list (через handle_request_table), |
|
|
|
|
но узел ещё не создан (входящий JOIN_GROUP). */ |
|
|
|
|
topo_group_conn_add_path(group, conn, conn->peer_node_id); |
|
|
|
|
{ struct TOPO_GROUP_NODE* peer_nq = topo_node_find_by_id(group, conn->peer_node_id); |
|
|
|
|
if (peer_nq) { |
|
|
|
|
peer_nq->connectivity.last_ping_time = get_time_tb(); |
|
|
|
|
if (conn->instance->topo_sqlite_db) topo_node_sqlite_nodeinfo_updated(conn->instance->topo_sqlite_db, conn->peer_node_id); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* дедуп: тот же conn стреляет ETCP_CONN_STATUS_UP дважды (UDP-линк, затем TCP-линк).
|
|
|
|
|
* Если conn уже в senders_list — повторно не обрабатываем, иначе active_conn_count |
|
|
|
|
* задваивается и переподключение никогда не стартует. */ |
|
|
|
|
* Если conn уже в senders_list — не слать повторно TABLE_REQ и не задваивать active_conn_count. */ |
|
|
|
|
{ struct ll_entry* se = group->senders_list ? group->senders_list->head : NULL; |
|
|
|
|
while (se) { |
|
|
|
|
if (((struct TOPO_GROUP_CONN_ITEM*)se->data)->conn == conn) { |
|
|
|
|
@ -480,15 +574,6 @@ void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
|
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
topo_recovery_cancel_for_node(group, conn->peer_node_id); |
|
|
|
|
|
|
|
|
|
struct TOPO_GROUP_NODE* peer_nq = topo_node_find_by_id(group, conn->peer_node_id); |
|
|
|
|
if (peer_nq) { |
|
|
|
|
peer_nq->connectivity.last_ping_time = get_time_tb(); |
|
|
|
|
peer_nq->conn_presence |= NCONN_DIRECT; |
|
|
|
|
peer_nq->conn_up |= NCONN_DIRECT; |
|
|
|
|
if (conn->instance->topo_sqlite_db) topo_node_sqlite_nodeinfo_updated(conn->instance->topo_sqlite_db, conn->peer_node_id); |
|
|
|
|
} |
|
|
|
|
topo_group_add_to_senders(group, conn); |
|
|
|
|
|
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "topo_group_new_conn: peer=%016llx group=%016llx type=%d ch=%s", (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->group_type, group->channel_id); |
|
|
|
|
@ -740,8 +825,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
|
|
|
|
|
if (nodeinfo1) { |
|
|
|
|
{ struct TOPO_NODE* stored = topo_node_registry_store(group->instance->topo_groups, new_ni); |
|
|
|
|
if (stored != new_ni) new_ni = stored; } |
|
|
|
|
new_ni->group_id = group->group_id; |
|
|
|
|
new_ni->flags = pkt->node.flags; |
|
|
|
|
nodeinfo1->node_id = new_ni ? new_ni->node_id : 0; |
|
|
|
|
nodeinfo1->subnets = new_subnets; |
|
|
|
|
} else { |
|
|
|
|
@ -751,8 +834,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
|
|
|
|
|
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); |
|
|
|
|
if (stored != new_ni) new_ni = stored; } |
|
|
|
|
new_ni->group_id = group->group_id; |
|
|
|
|
new_ni->flags = pkt->node.flags; |
|
|
|
|
nodeinfo1->node_id = new_ni ? new_ni->node_id : 0; |
|
|
|
|
nodeinfo1->subnets = new_subnets; |
|
|
|
|
nodeinfo1->connectivity.probe_status = PROBE_STATUS_NONE; |
|
|
|
|
@ -888,8 +969,9 @@ void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE*
|
|
|
|
|
u_free(p); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx ver=%d grp=%016llx to conn=%s", (unsigned long long)node->node_id, sni->ver, (unsigned long long)sni->group_id, conn->log_name); |
|
|
|
|
int ser_len = topo_node_serialize(sni, node, p + 2, max_sz - 2, cumulative_rtt); |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx ver=%d grp=%016llx to conn=%s", (unsigned long long)node->node_id, sni->ver, (unsigned long long)group->group_id, conn->log_name); |
|
|
|
|
uint8_t sflags = (group->group_type == TOPO_GROUP_TYPE_CHAT) ? 0 : TOPO_FLAG_SEND_SUBNETS; |
|
|
|
|
int ser_len = topo_node_serialize(sni, node, group->group_id, sflags, p + 2, max_sz - 2, cumulative_rtt); |
|
|
|
|
if (ser_len < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: serialize failed for node %016llx", (unsigned long long)node->node_id); u_free(p); return; } |
|
|
|
|
|
|
|
|
|
struct ll_entry* e = queue_entry_new(0); |
|
|
|
|
@ -941,10 +1023,19 @@ static void topo_group_handle_request_table(struct TOPO_GROUP* group, struct ETC
|
|
|
|
|
if (!group || !conn) return; |
|
|
|
|
topo_group_send_nodeinfo(group, group->local_node, conn, 0); |
|
|
|
|
topo_group_send_full_table(group, conn); |
|
|
|
|
topo_group_add_to_senders(group, conn); |
|
|
|
|
/* senders_list заполняется только через topo_group_new_conn (add + BGP);
|
|
|
|
|
здесь не добавляем, иначе new_conn (на JOIN_GROUP) упирается в дедуп и не шлёт TABLE_REQ обратно. */ |
|
|
|
|
topo_group_send_table_complete(group, conn); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { |
|
|
|
|
if (!group || !conn) return; |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_join_group: from %s grp=%016llx type=%d", conn->log_name, (unsigned long long)group->group_id, group->group_type); |
|
|
|
|
/* входящий запросил членство: для non-CHAT (UTUN) — добавить запросившего узла + инициировать BGP обратно */ |
|
|
|
|
if (group->group_type != TOPO_GROUP_TYPE_CHAT) |
|
|
|
|
topo_group_new_conn(group, conn); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* ── BGP node event callbacks ── */ |
|
|
|
|
|
|
|
|
|
void topo_group_add_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, void* arg) { |
|
|
|
|
|