diff --git a/src/chat/chat_event.h b/src/chat/chat_event.h index a8785964..30713d55 100644 --- a/src/chat/chat_event.h +++ b/src/chat/chat_event.h @@ -34,6 +34,7 @@ extern "C" { #define CHAT_EVT_SERVICE_STOPPED 13 /* data: none */ #define CHAT_EVT_ATTACHMENT_DOWNLOADED 14 /* [ch_id_len:1][ch_id:var][msg_id:8] */ #define CHAT_EVT_DOWNLOAD_PROGRESS 15 /* [ch_id_len:1][ch_id:var][msg_id:8][blocks_done:4][num_blocks:4] */ +#define CHAT_EVT_NODEINFO_UPDATED 16 /* [group_id:8][node_id:8][conn_presence:1][conn_up:1][best_rtt:2] */ typedef void (*chat_event_handler_fn)(int type, const uint8_t* data, int len); diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 311dcd65..ad284b75 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -600,61 +600,47 @@ struct chat_invite { int addr_count; }; +/* коллбэк invite-подключения: conn_mgr сам всё делает, здесь только лог результата. + реальный join-протокол запускается через cs_on_conn_up (глобальный callback), + сохранение в БД — через cs_handle_channel_info_resp после верификации */ +static void cs_invite_conn_cb(struct CONN_MGR_HANDLE* h, + uint64_t node_id, uint64_t group_id, + enum conn_mgr_event event, void* arg) { + (void)arg; + if (event == CONN_EVENT_JOIN) { + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: JOINED — connected to 0x%016llx, online now", + (unsigned long long)node_id); + struct UTUN_INSTANCE* inst = chat_core_get_inst(); + if (inst && inst->topo_groups) { + struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id); + if (g) { + struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, node_id); + if (nq) nq->handle = h; + } + } + } else if (event == CONN_EVENT_TIMEOUT) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "channel invite: TIMEOUT — could not connect to 0x%016llx", + (unsigned long long)node_id); + struct UTUN_INSTANCE* inst = chat_core_get_inst(); + if (inst && inst->topo_groups) { + struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id); + if (g) { + struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, node_id); + if (nq) nq->handle = NULL; + } + } + } +} + struct cm_invite_wrap { struct chat_invite inv; }; static void cm_invite_trampoline(void* arg) { struct cm_invite_wrap* w = (struct cm_invite_wrap*)arg; struct chat_invite* inv = &w->inv; struct UTUN_INSTANCE* inst = chat_core_get_inst(); - sqlite3* db = chat_core_get_db(); uint64_t node_id = inv->node_id; uint64_t channel_id = inv->channel_id; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start ch=%llu node=0x%016llx addrs=%d", - CS_ID, channel_id, node_id, inv->addr_count); - - /* save pubkey to nodes */ - sqlite3_stmt* st = NULL; - sqlite3_prepare_v2(db, "INSERT INTO nodes(node_id,x25519_pubkey,created_at) VALUES(?,?,?)" - " ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey," - " created_at=COALESCE(nodes.created_at, excluded.created_at)", -1, &st, NULL); - if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 3, 0); sqlite3_step(st); sqlite3_finalize(st); } - - /* save invite addresses (socket_id=0) */ - sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", -1, &st, NULL); - if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_step(st); sqlite3_finalize(st); } - sqlite3_prepare_v2(db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id) VALUES(?,?,?,?,?,0,?)", -1, &st, NULL); - if (st) { - const uint8_t* ap = inv->addrs_data; - for (int i = 0; i < inv->addr_count; i++) { - uint8_t family = *ap++; ap++; /* skip sock_id */ - uint8_t proto = *ap++; /* read proto */ - if (family == 4) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_int(st, 2, 4); sqlite3_bind_int(st, 3, proto); - sqlite3_bind_blob(st, 4, ap, 4, SQLITE_STATIC); ap += 4; - uint16_t port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; - sqlite3_bind_int(st, 5, (int)port); sqlite3_bind_int(st, 6, 1); - sqlite3_step(st); sqlite3_reset(st); - } else if (family == 6) { - sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_int(st, 2, 6); sqlite3_bind_int(st, 3, proto); - sqlite3_bind_blob(st, 4, ap, 16, SQLITE_STATIC); ap += 16; - uint16_t port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; - sqlite3_bind_int(st, 5, (int)port); sqlite3_bind_int(st, 6, 1); - sqlite3_step(st); sqlite3_reset(st); - } else { ap += 18; } - } - sqlite3_finalize(st); - } - - char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)channel_id); - - /* TOPO_GROUP (lightweight) for conn_mgr_connect_from_invite below. - DB channel + sync created later in cs_handle_channel_info_resp. */ - if (inst->topo_groups) { - if (!topo_groups_find(inst->topo_groups, channel_id)) - topo_groups_create_group(inst->topo_groups, channel_id, TOPO_GROUP_TYPE_CHAT, ch_str); - } - /* build TOPO_NODE from invite data */ struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); if (!ni) { u_free(w); return; } @@ -664,8 +650,8 @@ static void cm_invite_trampoline(void* arg) { uint8_t family = *ap++; uint8_t sock_id = *ap++; uint8_t proto = *ap++; uint16_t port = 0; if (family == 4) { port = ((uint16_t)ap[4] << 8) | ap[5]; } - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: invite addr[%d] family=%d proto=%d sock=%d port=%d", - CS_ID, i, family, proto, sock_id, port); + else if (family == 6) { port = ((uint16_t)ap[16] << 8) | ap[17]; } + const char* proto_str = proto == 1 ? "UDP" : proto == 2 ? "TCP" : proto == 3 ? "UDP+TCP" : "?"; if (family == 4) { struct TOPO_ADDR4* a4 = u_calloc(1, sizeof(*a4)); if (!a4) break; @@ -683,32 +669,20 @@ static void cm_invite_trampoline(void* arg) { } else { ap += 18; } } - struct ETCP_CONN* existing = instance_find_conn(inst, node_id); - if (!existing) existing = cs_find_conn_by_pubkey(inst, inv->pubkey, &node_id); - - if (existing) { - if (existing->links_up) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start — already connected peer=0x%016llx, joining ch=%llu", - CS_ID, (unsigned long long)node_id, (unsigned long long)channel_id); - if (node_id != existing->peer_node_id) - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: invite node_id mismatch: derived=%016llx ETCP=%016llx — using ETCP", - CS_ID, (unsigned long long)node_id, (unsigned long long)existing->peer_node_id); - g_cs->pending_invite_ch_id = channel_id; - g_cs->pending_invite_node_id = existing->peer_node_id; - cs_start_channel_join(g_cs, existing->peer_node_id); - cs_on_peer_status_changed(existing->peer_node_id, 1); - } else { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite conn in progress node=0x%016llx, waiting UP", CS_ID, (unsigned long long)node_id); - } - } else if (inst->topo_groups && (ni->v4_addrs || ni->v6_addrs)) { - struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id); - if (g && g->conn_mgr) { - conn_mgr_open_invite(inst, channel_id, ni, node_id, NULL, NULL, NULL); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite connect posted node=0x%016llx", CS_ID, (unsigned long long)node_id); - } - } else { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: invite no addrs topo=%p v4=%p v6=%p", CS_ID, - (void*)inst->topo_groups, (void*)ni->v4_addrs, (void*)ni->v6_addrs); + if (!ni->v4_addrs && !ni->v6_addrs) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "invite: no valid addrs for node=0x%016llx", (unsigned long long)node_id); + while (ni->v4_addrs) { struct TOPO_ADDR4* n = ni->v4_addrs->next; u_free(ni->v4_addrs); ni->v4_addrs = n; } + while (ni->v6_addrs) { struct TOPO_ADDR6* n = ni->v6_addrs->next; u_free(ni->v6_addrs); ni->v6_addrs = n; } + u_free(ni); u_free(w); return; + } + + /* conn_mgr_open_invite сам создаст TOPO_GROUP, найдёт/создаст соединение, запустит INVITE_INFO. + DB channel + sync создаются позже в cs_handle_channel_info_resp после верификации */ + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "channel invite: connecting to 0x%016llx via %d addresses", + (unsigned long long)node_id, inv->addr_count); + int r = conn_mgr_open_invite(inst, channel_id, ni, node_id, cs_invite_conn_cb, NULL, NULL); + if (r < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "channel invite: FAILED — conn_mgr_open_invite error=%d", r); } /* free caller-owned ni */ @@ -720,7 +694,8 @@ static void cm_invite_trampoline(void* arg) { void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, - const uint8_t* addrs_data, int addr_count) { + const uint8_t* addrs_data, int addr_count, + int addrs_data_len) { if (!g_cs || !g_cs->inst || !g_cs->inst->ua) { int r = -7; #ifdef __ANDROID__ @@ -741,7 +716,7 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, inv->node_id = node_id; memcpy(inv->pubkey, pubkey_bin, 32); - size_t addrs_sz = (size_t)addr_count * 9; + size_t addrs_sz = addrs_data_len > 0 ? (size_t)addrs_data_len : (size_t)addr_count * 9; inv->addrs_data = u_malloc(addrs_sz); if (!inv->addrs_data) { u_free(inv); return; } memcpy(inv->addrs_data, addrs_data, addrs_sz); diff --git a/src/chat/chat_sync.h b/src/chat/chat_sync.h index fa5550b0..3b620d1c 100644 --- a/src/chat/chat_sync.h +++ b/src/chat/chat_sync.h @@ -75,7 +75,8 @@ void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); /* Подключиться к пиру по данным invite-ссылки (вызывается из uasync-потока) */ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, - const uint8_t* addrs_data, int addr_count); + const uint8_t* addrs_data, int addr_count, + int addrs_data_len); /* Присоединиться к каналу через уже подключённый узел (uasync-поток). Соединение с target_node_id должно быть установлено (links_up && initialized). diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index a4abad7c..46160273 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/src/routing_layer/conn_mgr_core.c @@ -206,6 +206,9 @@ void cm_update_nodeinfo(struct CONN_MGR_ENTRY* entry) { nq->conn_mgr_type = entry->conn_type; memcpy(nq->conn_mgr_intermediaries, entry->intermediaries, sizeof(entry->intermediaries)); nq->conn_mgr_intermediariy_count = entry->intermediariy_count; + if (entry->conn_type == CONN_TYPE_INDIRECT) { nq->conn_presence |= NCONN_INDIRECT; nq->conn_up |= NCONN_INDIRECT; } + else { nq->conn_presence &= ~NCONN_INDIRECT; nq->conn_up &= ~NCONN_INDIRECT; } + topo_fire_nodeinfo_cbk(entry->mgr->instance, entry->mgr->group, nq); } void cm_cleanup_db_node(struct CONN_MGR_ENTRY* entry) { @@ -235,14 +238,16 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void switch (ncd_ev) { case NCD_EVENT_UP: if (entry->state == CONN_MGR_STATE_CONNECTED) return; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: NCD UP for 0x%016llx", (unsigned long long)entry->node_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP handshake OK with 0x%016llx — connection ESTABLISHED", + (unsigned long long)entry->node_id); if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; } entry->main_connect_state = CM_TRY_OK; entry->conn_type = CONN_TYPE_DIRECT; entry->state = CONN_MGR_STATE_CONNECTED; cm_update_nodeinfo(entry); cm_deliver_event(entry, CONN_EVENT_UP); break; case NCD_EVENT_TIMEOUT: - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: NCD TIMEOUT for 0x%016llx", (unsigned long long)entry->node_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP connect TIMEOUT to 0x%016llx — no response, trying fallback", + (unsigned long long)entry->node_id); entry->ncd_handle = NULL; if (entry->db_loaded) { entry->main_connect_state = CM_TRY_FAILED; @@ -261,7 +266,8 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void if (entry->local_scan_state == CM_TRY_FAILED) cm_deliver_event(entry, CONN_EVENT_TIMEOUT); break; case NCD_EVENT_DOWN: - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: NCD DOWN for 0x%016llx", (unsigned long long)entry->node_id); + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "link: ETCP link DROPPED to 0x%016llx", + (unsigned long long)entry->node_id); entry->ncd_handle = NULL; cm_deliver_event(entry, CONN_EVENT_DOWN); break; @@ -275,7 +281,8 @@ 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 group=0x%016llx", (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: ETCP handshake OK with 0x%016llx — sending membership check", + (unsigned long long)inv->node_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; @@ -284,15 +291,15 @@ void cm_invite_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_e qe->dgram = u_malloc(sizeof(req)); if (!qe->dgram) { queue_entry_free(qe); cm_invite_fail(inv); return; } memcpy(qe->dgram, &req, sizeof(req)); qe->len = sizeof(req); etcp_send(node_conn_direct_get_conn(ncd_h), qe); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_REQ to 0x%016llx group=0x%016llx", - (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id); } break; case NCD_EVENT_TIMEOUT: - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite connect timeout node=0x%016llx", (unsigned long long)inv->node_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: ETCP connect TIMEOUT to 0x%016llx — no response from peer", + (unsigned long long)inv->node_id); cm_invite_fail(inv); break; case NCD_EVENT_DOWN: - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite conn down node=0x%016llx", (unsigned long long)inv->node_id); + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: ETCP link DROPPED to 0x%016llx while connecting", + (unsigned long long)inv->node_id); if (inv->state == CM_INVITE_CONNECTING) cm_invite_fail(inv); break; } @@ -711,9 +718,8 @@ void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry /* ═══════ invite обработчики ═══════ */ void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn, const uint8_t* data, size_t len) { - if (len < sizeof(struct CM_INVITE_REQ)) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ too short (%zu)", len); return; } + if (len < sizeof(struct CM_INVITE_REQ)) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: INVITE_INFO_REQ too short (%zu)", len); return; } struct CM_INVITE_REQ* req = (struct CM_INVITE_REQ*)data; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_REQ from 0x%016llx group=0x%016llx", (unsigned long long)conn->peer_node_id, (unsigned long long)req->group_id); struct TOPO_GROUP* grp = topo_groups_find(mgr->instance->topo_groups, req->group_id); struct TOPO_GROUP* sg = grp ? grp : mgr->group; struct TOPO_NODE* ln = sg->local_node ? topo_node_registry_find(sg->instance->topo_groups, sg->local_node->node_id) : NULL; @@ -725,19 +731,21 @@ void cm_handle_invite_info_req(struct CONN_MGR* mgr, struct ETCP_CONN* conn, con if (nl) memcpy(resp->node_name, nm, nl); struct ll_entry* qe = queue_entry_new(0); if (qe) { qe->dgram=(uint8_t*)resp; qe->len=(uint16_t)rs; etcp_send(conn, qe); } else u_free(resp); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: sent INVITE_INFO_RESP to 0x%016llx name_len=%zu", (unsigned long long)conn->peer_node_id, nl); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: membership check RECEIVED from 0x%016llx — responding is_member=yes name=\"%.*s\"", + (unsigned long long)conn->peer_node_id, (int)nl, nm ? nm : ""); } void cm_handle_invite_info_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) { - if (len < sizeof(struct CM_INVITE_RESP)) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP too short (%zu)", len); return; } + if (len < sizeof(struct CM_INVITE_RESP)) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: INVITE_INFO_RESP too short (%zu)", len); return; } struct CM_INVITE_RESP* resp = (struct CM_INVITE_RESP*)data; - if (len < sizeof(struct CM_INVITE_RESP) + (size_t)resp->node_name_len) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP truncated"); return; } + if (len < sizeof(struct CM_INVITE_RESP) + (size_t)resp->node_name_len) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: INVITE_INFO_RESP truncated"); return; } struct cm_invite_pending* inv = NULL; for (struct cm_invite_pending* p = mgr->invite_list; p; p = p->next) if (p->state == CM_INVITE_WAIT_INFO) { inv = p; break; } - if (!inv) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP — no pending invite in WAIT_INFO"); return; } - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INVITE_INFO_RESP node=0x%016llx group=0x%016llx member=%d name_len=%u", - (unsigned long long)inv->node_id, (unsigned long long)resp->group_id, resp->is_member, resp->node_name_len); - if (!resp->is_member) { cm_invite_fail(inv); return; } + if (!inv) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: INVITE_INFO_RESP — no pending invite in WAIT_INFO state"); return; } + if (!resp->is_member) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "invite: membership DENIED by 0x%016llx — not in channel", (unsigned long long)inv->node_id); + cm_invite_fail(inv); return; + } if (inv->temp_nq) { struct TOPO_NODE* ni2 = topo_node_registry_find(inv->mgr->instance->topo_groups, inv->temp_nq->node_id); if (ni2) { memcpy(ni2->ed25519_public_key, resp->ed25519_pubkey, SC_PUBKEY_SIZE); @@ -746,16 +754,16 @@ 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 : ""); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: membership CONFIRMED by 0x%016llx name=\"%.*s\" — joined to channel, handle registered", + (unsigned long long)inv->node_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); } void cm_invite_overall_timeout(void* arg) { struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; inv->overall_timer = NULL; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: invite overall timeout node=0x%016llx", (unsigned long long)inv->node_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: overall TIMEOUT to 0x%016llx — no response within %ds", + (unsigned long long)inv->node_id, CM_INVITE_DEFAULT_TIMEOUT_MS/1000); cm_invite_fail(inv); } @@ -804,22 +812,21 @@ 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 group=0x%016llx", (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: TCP handshake OK with 0x%016llx — sending membership check", + (unsigned long long)inv->node_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; struct ll_entry* qe2 = queue_entry_new(0); if (!qe2) { cm_invite_fail(inv); return; } qe2->dgram=u_malloc(sizeof(req)); if (!qe2->dgram) { queue_entry_free(qe2); cm_invite_fail(inv); return; } memcpy(qe2->dgram,&req,sizeof(req)); qe2->len=sizeof(req); etcp_send(etcp, qe2); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP invite sent INVITE_INFO_REQ to 0x%016llx group=0x%016llx", - (unsigned long long)inv->node_id, (unsigned long long)inv->mgr->group->group_id); } } static void cm_tcp_close_cb(struct stcp_link* link, int err, void* arg) { struct cm_invite_pending* inv = (struct cm_invite_pending*)arg; if (!inv) return; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: TCP invite closed node=0x%016llx err=%d", (unsigned long long)inv->node_id, err); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: TCP link CLOSED to 0x%016llx err=%d", (unsigned long long)inv->node_id, err); inv->tcp_link = NULL; } @@ -833,9 +840,8 @@ static int cm_invite_tcp_connect(struct cm_invite_pending* inv, const struct TOP .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); - 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); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: TCP connecting to %d.%d.%d.%d:%d", + a->addr[0], a->addr[1], a->addr[2], a->addr[3], (int)a->port); return 1; } return 0; } @@ -852,8 +858,8 @@ static int cm_invite_tcp_connect(struct cm_invite_pending* inv, const struct TOP .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); - 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); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: TCP6 connecting to port %d", + (int)a6->port); return 1; } return 0; } @@ -929,7 +935,7 @@ int conn_mgr_open_invite(struct UTUN_INSTANCE* inst, inv->overall_timer = uasync_set_timeout(inst->ua, CM_INVITE_DEFAULT_TIMEOUT_MS*10, inv, cm_invite_overall_timeout, "conn_mgr_invite"); { - char addrs[256] = ""; int aoff = 0; + char addrs[512] = ""; 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, @@ -937,8 +943,8 @@ int conn_mgr_open_invite(struct UTUN_INSTANCE* inst, 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); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "invite: starting to 0x%016llx via%s udp=%d tcp=%d timeout=%ds", + (unsigned long long)nid, addrs, udp, tcp, CM_INVITE_DEFAULT_TIMEOUT_MS/1000); } return 0; } diff --git a/src/routing_layer/conn_mgr_indirect.c b/src/routing_layer/conn_mgr_indirect.c index c72b7c8c..c43b08f1 100644 --- a/src/routing_layer/conn_mgr_indirect.c +++ b/src/routing_layer/conn_mgr_indirect.c @@ -98,7 +98,7 @@ void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, s uint8_t need=0; uint64_t now=get_time_tb(); for(uint8_t i=0;icached.my_count&&i<4;i++){ struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,ep->cached.my_candidates[i].node_id); - if(!nq||!cm_is_rtt_fresh(nq,now)){need++;route_connectivity_probe_node(mgr->instance,nq);} + if(!nq||!cm_is_rtt_fresh(nq,now)){need++;route_connectivity_probe_node(mgr->instance,mgr->group,nq);} } if(!need){ep->probed=1;cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}u_free(ep);return;} ep->timer=uasync_set_timeout(mgr->instance->ua,2000,ep,cm_exchange_probe_retry_cb,"cm_exch_probe"); @@ -126,7 +126,7 @@ void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CM_EXCHANGE_RE uint64_t now=get_time_tb(); for(uint8_t i=0;icandidate_count&&i<4;i++){ struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,req->candidates[i].node_id); - if(nq&&cm_is_rtt_fresh(nq,now))continue; route_connectivity_probe_node(mgr->instance,nq); + if(nq&&cm_is_rtt_fresh(nq,now))continue; route_connectivity_probe_node(mgr->instance,mgr->group,nq); } struct CM_EXCHANGE_RESP resp; memset(&resp,0,sizeof(resp)); resp.cmd=ETCP_RT_ID_CONN_MGR; resp.subcmd=CM_SUBCMD_INTERM_EXCHANGE_RESP; resp.request_id=req->request_id; diff --git a/src/routing_layer/conn_mgr_monitor.c b/src/routing_layer/conn_mgr_monitor.c index 0df701d4..81597eb4 100644 --- a/src/routing_layer/conn_mgr_monitor.c +++ b/src/routing_layer/conn_mgr_monitor.c @@ -169,7 +169,7 @@ void cm_candidate_ping_timer_cb(void* arg) { struct ETCP_CONN* dc = topo_group_find_conn_for_node(mgr->group, nid); struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, nq->node_id); if (ni && dc) etcp_send_ping(mgr->instance, ni->public_key, &dc->links->remote_addr, CONN_PROBE_TIMEOUT_MS, NULL, NULL, NULL, 0); - } else route_connectivity_probe_node(mgr->instance, nq); + } else route_connectivity_probe_node(mgr->instance, mgr->group, nq); } 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"); } diff --git a/src/routing_layer/route_connectivity.c b/src/routing_layer/route_connectivity.c index 9243a9a4..ff1c6e80 100644 --- a/src/routing_layer/route_connectivity.c +++ b/src/routing_layer/route_connectivity.c @@ -23,6 +23,7 @@ struct conn_probe_ctx { struct conn_probe_ctx* next; /* linked list in nq->connectivity.probe_list */ struct UTUN_INSTANCE* instance; + struct TOPO_GROUP* group; struct TOPO_GROUP_NODE* nq; uint8_t addr_type; // ADDR_TYPE_* struct sockaddr_storage target_addr; @@ -283,6 +284,7 @@ static void conn_probe_finish(struct conn_probe_ctx* ctx, int ok) { c->interface_status, c->nat_status, c->real_status); if (ctx->instance && ctx->instance->topo_sqlite_db) topo_node_sqlite_nodeinfo_updated(ctx->instance->topo_sqlite_db, ctx->nq->node_id); + topo_fire_nodeinfo_cbk(ctx->instance, ctx->group, ctx->nq); } } @@ -291,7 +293,7 @@ static void conn_probe_finish(struct conn_probe_ctx* ctx, int ok) { // ---- public API ---- -void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_GROUP_NODE* nq) { +void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq) { if (!instance || !nq) return; struct TOPO_NODE* ni = topo_node_registry_find(instance->topo_groups, nq->node_id); if (!ni) return; @@ -332,7 +334,7 @@ void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_G if (cand_count == 0) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "probe: no local sockets for target %s:%u type=%d", ip_to_str(&ip, AF_INET).str, port, a->type); continue; } struct conn_probe_ctx* ctx = u_calloc(1, sizeof(struct conn_probe_ctx)); if (!ctx) continue; - ctx->instance = instance; ctx->nq = nq; ctx->addr_type = a->type; ctx->target_addr = target; + ctx->instance = instance; ctx->group = group; ctx->nq = nq; ctx->addr_type = a->type; ctx->target_addr = target; memcpy(ctx->peer_pubkey, ni->public_key, SC_PUBKEY_SIZE); memcpy(ctx->candidate_sockets, candidates, cand_count * sizeof(struct ETCP_SOCKET*)); ctx->candidate_count = cand_count; ctx->candidate_index = 0; ctx->count_total = CONN_PROBE_COUNT; ctx->timeout_ms = CONN_PROBE_TIMEOUT_MS; ctx->best_across_sockets = 65535; diff --git a/src/routing_layer/route_connectivity.h b/src/routing_layer/route_connectivity.h index 6bcb9372..9e47b2ba 100644 --- a/src/routing_layer/route_connectivity.h +++ b/src/routing_layer/route_connectivity.h @@ -10,9 +10,10 @@ extern "C" { #include "topo_node.h" struct UTUN_INSTANCE; +struct TOPO_GROUP; // Запускает зондирование связности ко всем адресам удалённого узла -void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, +void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq); // Отменяет все pending пробы для узла (при удалении / withdraw) diff --git a/src/routing_layer/route_lib.c b/src/routing_layer/route_lib.c index a97c99d2..ce2f88b2 100644 --- a/src/routing_layer/route_lib.c +++ b/src/routing_layer/route_lib.c @@ -107,13 +107,14 @@ static struct ROUTE_ENTRY* binary_search_lpm(struct ROUTE_TABLE *table, uint32_t // Основные функции // ============================================================================ -struct ROUTE_TABLE *route_table_create(void) { +struct ROUTE_TABLE *route_table_create(uint64_t my_node_id) { struct ROUTE_TABLE *table = u_calloc(1, sizeof(struct ROUTE_TABLE)); if (!table) { DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "alloc failed"); return NULL; } + table->my_node_id = my_node_id; table->entries = u_calloc(INITIAL_CAPACITY, sizeof(struct ROUTE_ENTRY)); if (!table->entries) { DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "alloc entries failed"); @@ -180,7 +181,7 @@ bool route_insert(struct ROUTE_TABLE *table, struct TOPO_GROUP_NODE *node) { struct ROUTE_ENTRY* e = &table->entries[pos]; e->network = network; e->prefix_length = s->prefix_length; e->v_node_info = node; table->count++; - if (node->hop_count == 0) table->stats.local_routes++; else table->stats.learned_routes++; + if (node->node_id == table->my_node_id) table->stats.local_routes++; else table->stats.learned_routes++; } DEBUG_INFO(DEBUG_CATEGORY_ROUTING, "added %zu route(s) for node=%016llx", @@ -196,7 +197,7 @@ void route_delete(struct ROUTE_TABLE *table, struct TOPO_GROUP_NODE *node) { if (table->entries[i].v_node_info == node) { if (i < table->count - 1) memmove(&table->entries[i], &table->entries[i + 1], (table->count - i - 1) * sizeof(struct ROUTE_ENTRY)); table->count--; - if (node->hop_count == 0) table->stats.local_routes--; else table->stats.learned_routes--; + if (node->node_id == table->my_node_id) table->stats.local_routes--; else table->stats.learned_routes--; removed++; continue; } i++; @@ -236,7 +237,7 @@ void route_table_print(const struct ROUTE_TABLE *table) { const struct ROUTE_ENTRY *entry = &table->entries[i]; DEBUG_INFO(DEBUG_CATEGORY_ROUTING, " %zu: %s/%d", i + 1, ip_to_string(entry->network).a, entry->prefix_length); - if (entry->v_node_info && entry->v_node_info->hop_count != 0) { + if (entry->v_node_info && entry->v_node_info->node_id != table->my_node_id) { DEBUG_INFO(DEBUG_CATEGORY_ROUTING, " v_node_info=%p", (void*)entry->v_node_info); } else { DEBUG_INFO(DEBUG_CATEGORY_ROUTING, " LOCAL"); diff --git a/src/routing_layer/route_lib.h b/src/routing_layer/route_lib.h index 657e7038..bbcfed42 100644 --- a/src/routing_layer/route_lib.h +++ b/src/routing_layer/route_lib.h @@ -47,6 +47,7 @@ struct ROUTE_TABLE { struct ROUTE_ENTRY *entries; /**< Массив записей маршрутов */ size_t count; /**< Текущее количество записей */ size_t capacity; /**< Максимальная емкость таблицы */ + uint64_t my_node_id; /**< ID локального узла */ uint32_t *dynamic_subnets; /**< Динамические подсети (массив пар: сеть, префикс) */ size_t dynamic_subnet_count; /**< Количество динамических подсетей */ uint32_t *local_subnets; /**< Локальные подсети (массив пар: сеть, префикс) */ @@ -64,7 +65,7 @@ struct ROUTE_TABLE { * * @return Указатель на созданную таблицу или NULL в случае ошибки */ -struct ROUTE_TABLE *route_table_create(void); +struct ROUTE_TABLE *route_table_create(uint64_t my_node_id); /** * @brief Уничтожает таблицу маршрутизации и освобождает ресурсы diff --git a/src/routing_layer/routing.c b/src/routing_layer/routing.c index 5211d609..a9be50e2 100644 --- a/src/routing_layer/routing.c +++ b/src/routing_layer/routing.c @@ -140,7 +140,7 @@ void route_pkt(struct UTUN_INSTANCE* instance, struct ll_entry* entry, uint64_t struct TOPO_GROUP_NODE* nq = route->v_node_info; struct TOPO_NODE* rni = nq ? topo_node_registry_find(instance->topo_groups, nq->node_id) : NULL; - if (!nq || !rni || nq->hop_count == 0) { + if (!nq || !rni || nq->node_id == instance->node_id) { DEBUG_TRACE(DEBUG_CATEGORY_ROUTING, "Local route to %s", ip_to_str(&addr, AF_INET).str); } else { DEBUG_TRACE(DEBUG_CATEGORY_ROUTING, "route_pkt: sending %zu bytes to node %016llx dst=%s", @@ -233,7 +233,7 @@ int routing_create(struct UTUN_INSTANCE* instance) { } // Create route table - instance->rt = route_table_create(); + instance->rt = route_table_create(instance->node_id); if (!instance->rt) { DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "failed to create route table for node %016llx", (unsigned long long)instance->node_id); diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 9508bb03..a474dea8 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -201,6 +201,10 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg 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) { + struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, conn->peer_node_id); + if (nq) { nq->conn_presence &= ~NCONN_DIRECT; nq->conn_up &= ~NCONN_DIRECT; } + } } } @@ -408,6 +412,12 @@ 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_DEBUG, "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); @@ -434,21 +444,34 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { while (node_entry) { struct ll_entry* next = node_entry->next; struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)node_entry; - 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, 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); - 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); - struct ll_entry* entry = node_entry; - if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); } - nodes_removed++; - etcp_router_conn_close_all_for_node(group->instance, group->group_id, key); - uint64_t wd_src = (key == conn->peer_node_id) ? conn->instance->node_id : conn->peer_node_id; - topo_group_broadcast_withdraw(group, key, wd_src, NULL); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Removed node %016llx after link down", (unsigned long long)key); + { + uint64_t next_hop = nq->node_id; + struct ll_entry* pe = nq->paths ? nq->paths->head : NULL; + 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++; } + 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); + if (conn->instance && conn->instance->control_srv) control_server_notify_node_removed(conn->instance->control_srv, key); + 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); + struct ll_entry* entry = node_entry; + if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); } + nodes_removed++; + etcp_router_conn_close_all_for_node(group->instance, group->group_id, key); + uint64_t wd_src = (key == conn->peer_node_id) ? conn->instance->node_id : conn->peer_node_id; + topo_group_broadcast_withdraw(group, key, wd_src, NULL); + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "Removed node %016llx after link down", (unsigned long long)key); + } else { + nq->conn_up &= ~NCONN_DIRECT; + { int has_direct = 0; struct ll_entry* pe2 = nq->paths ? nq->paths->head : NULL; + while (pe2) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)pe2; if (p->conn->peer_node_id == nq->node_id) { has_direct = 1; break; } pe2 = pe2->next; } + if (!has_direct) nq->conn_presence &= ~NCONN_DIRECT; } + topo_fire_nodeinfo_cbk(conn->instance, group, nq); + } } node_entry = next; } @@ -579,14 +602,15 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from memcpy(hop_list, hop_src, ni->hop_count * 8); hop_list[ni->hop_count] = from->peer_node_id; topo_group_remove_path(nodeinfo1, from); - topo_group_add_path(nodeinfo1, from, hop_list, new_hops, nodeinfo1->cumulative_rtt); + { uint8_t _hc; uint16_t _rtt; topo_node_best_hop_list(nodeinfo1, &_hc, &_rtt); + topo_group_add_path(nodeinfo1, from, hop_list, new_hops, _rtt); } } return 0; } DEBUG_INFO(DEBUG_CATEGORY_BGP, "NODEINFO: node=%016llx ver=%d v4s=%d v4a=%d v6s=%d v6a=%d hops=%d from=%s", - (unsigned long long)node_id, new_ver, ni->local_v4_sockets, ni->local_v4_addrs, - ni->local_v6_sockets, ni->local_v6_addrs, ni->hop_count, from->log_name); + (unsigned long long)node_id, new_ver, ni->local_v4_sockets, ni->local_v4_addrs, + ni->local_v6_sockets, ni->local_v6_addrs, ni->hop_count, from->log_name); struct ll_queue* paths = NULL; int is_new_node = 0; @@ -606,6 +630,10 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from if (nodeinfo1) { queue_remove_data(group->nodes, &nodeinfo1->ll); queue_free(paths); queue_entry_free(&nodeinfo1->ll); } return -1; } + { int v4c = topo_list_count((struct _topo_head*)new_ni->v4_addrs); + int v6c = topo_list_count((struct _topo_head*)new_ni->v6_addrs); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "NODEINFO deser: nid=%016llx v4a=%d v6a=%d from=%s", + (unsigned long long)node_id, v4c, v6c, from->log_name); } DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "NODEINFO deserialized: nid=%016llx new_ver=%d flags=0x%02x ed=%016llx", (unsigned long long)node_id, new_ni->ver, new_ni->flags, *(uint64_t*)new_ni->ed25519_public_key); @@ -637,9 +665,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from new_ni->flags = pkt->node.flags; nodeinfo1->node_id = new_ni ? new_ni->node_id : 0; nodeinfo1->subnets = new_subnets; - nodeinfo1->cumulative_rtt = incoming_cumulative_rtt; - 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_hop_list); u_free(new_subnets); return -1; } @@ -653,8 +678,6 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from new_ni->flags = pkt->node.flags; nodeinfo1->node_id = new_ni ? new_ni->node_id : 0; nodeinfo1->subnets = new_subnets; - nodeinfo1->cumulative_rtt = incoming_cumulative_rtt; - 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; nodeinfo1->connectivity.nat_status = PROBE_RESULT_UNKNOWN; @@ -666,16 +689,17 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from nodeinfo1->paths = paths; nodeinfo1->last_ver = new_ver; uint64_t hop_list[MAX_HOPS]; - if (nodeinfo1->hop_count > 0) memcpy(hop_list, nodeinfo1->hop_list, nodeinfo1->hop_count * 8); - hop_list[nodeinfo1->hop_count] = from->peer_node_id; - nodeinfo1->hop_count++; + if (new_hop_count > 0) memcpy(hop_list, new_hop_list, new_hop_count * 8); + hop_list[new_hop_count] = from->peer_node_id; + uint8_t extended_count = new_hop_count + 1; - u_free(nodeinfo1->hop_list); - nodeinfo1->hop_list = u_malloc(nodeinfo1->hop_count * 8); - memcpy(nodeinfo1->hop_list, hop_list, nodeinfo1->hop_count * 8); + u_free(new_hop_list); topo_group_remove_path_by_hop(nodeinfo1, from->peer_node_id); - topo_group_add_path(nodeinfo1, from, hop_list, nodeinfo1->hop_count, incoming_cumulative_rtt); + topo_group_add_path(nodeinfo1, from, hop_list, extended_count, incoming_cumulative_rtt); + + nodeinfo1->conn_presence |= NCONN_BGP; nodeinfo1->conn_up |= NCONN_BGP; + if (from->peer_node_id == node_id) { nodeinfo1->conn_presence |= NCONN_DIRECT; nodeinfo1->conn_up |= NCONN_DIRECT; } if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) route_insert(group->instance->rt, nodeinfo1); if (node_id != group->instance->node_id) { @@ -700,8 +724,9 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from struct topo_node_cbk_entry* c = group->node_cbks; while (c) { c->fn(group, node_id, ev, c->arg); c = c->next; } } + topo_fire_nodeinfo_cbk(group->instance, group, nodeinfo1); - int hop_count = nodeinfo1->hop_count; + int hop_count = extended_count; struct ll_entry* se = group->senders_list ? group->senders_list->head : NULL; while (se) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; @@ -725,7 +750,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from struct TOPO_NODE* tni = topo_node_registry_find(group->instance->topo_groups, nodeinfo1->node_id); int v4a = tni ? topo_list_count((struct _topo_head*)tni->v4_addrs) : 0; int v6a = tni ? topo_list_count((struct _topo_head*)tni->v6_addrs) : 0; - if (v4a > 0 || v6a > 0) route_connectivity_probe_node(group->instance, nodeinfo1); + if (v4a > 0 || v6a > 0) route_connectivity_probe_node(group->instance, group, nodeinfo1); } return 0; @@ -740,6 +765,8 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send if (!nq) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "node not found"); return 0; } int ret = topo_group_remove_path_by_hop(nq, wd_source); if (ret == 1) { + nq->conn_presence = 0; nq->conn_up = 0; + topo_fire_nodeinfo_cbk(group->instance, group, nq); 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); route_connectivity_cancel_node(group->instance, nq); @@ -811,9 +838,9 @@ static void topo_group_send_full_table(struct TOPO_GROUP* group, struct ETCP_CON struct ll_entry* e = group->nodes ? group->nodes->head : NULL; while (e) { struct TOPO_GROUP_NODE* nq = (struct TOPO_GROUP_NODE*)e; - if (nq->hop_count == 0) { e = e->next; continue; } + if (nq->node_id == group->instance->node_id) { e = e->next; continue; } if (!nq->paths) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "node has no paths"); e = e->next; continue; } - if (topo_group_should_send_to(nq, target)) topo_group_send_nodeinfo(group, nq, conn, (uint16_t)((uint32_t)conn->rtt_last + nq->cumulative_rtt)); + if (topo_group_should_send_to(nq, target)) { uint8_t _hc; uint16_t _rtt; topo_node_best_hop_list(nq, &_hc, &_rtt); topo_group_send_nodeinfo(group, nq, conn, (uint16_t)((uint32_t)conn->rtt_last + _rtt)); } e = e->next; } } @@ -850,6 +877,36 @@ void topo_group_remove_node_cbk(struct TOPO_GROUP* group, topo_node_event_fn fn, } } +/* ── Global nodeinfo callbacks (instance-level, any node change) ── */ + +void utun_add_nodeinfo_cbk(struct UTUN_INSTANCE* instance, nodeinfo_cbk_fn fn, void* arg) { + if (!instance || !fn) return; + struct nodeinfo_cbk_entry* e = u_malloc(sizeof(*e)); + if (!e) return; + e->fn = fn; e->arg = arg; + e->next = instance->nodeinfo_cbks; + instance->nodeinfo_cbks = e; +} + +void utun_remove_nodeinfo_cbk(struct UTUN_INSTANCE* instance, nodeinfo_cbk_fn fn, void* arg) { + if (!instance || !fn) return; + struct nodeinfo_cbk_entry** p = &instance->nodeinfo_cbks; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct nodeinfo_cbk_entry* r = *p; + *p = r->next; u_free(r); + return; + } + p = &(*p)->next; + } +} + +void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node) { + if (!instance || !group || !node) return; + struct nodeinfo_cbk_entry* c = instance->nodeinfo_cbks; + while (c) { c->fn(group, node, c->arg); c = c->next; } +} + void topo_group_send_withdraw(struct TOPO_GROUP* group, uint64_t node_id) { if (!group) return; DEBUG_INFO(DEBUG_CATEGORY_BGP, "node=%016llx", (unsigned long long)node_id); diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 16fc9dfd..5678268d 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -71,6 +71,20 @@ struct topo_node_cbk_entry { struct topo_node_cbk_entry* next; }; +/* ── Global nodeinfo callbacks (one per instance, any node change) ── */ + +typedef void (*nodeinfo_cbk_fn)(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node, void* arg); + +struct nodeinfo_cbk_entry { + nodeinfo_cbk_fn fn; + void* arg; + struct nodeinfo_cbk_entry* next; +}; + +void utun_add_nodeinfo_cbk(struct UTUN_INSTANCE* instance, nodeinfo_cbk_fn fn, void* arg); +void utun_remove_nodeinfo_cbk(struct UTUN_INSTANCE* instance, nodeinfo_cbk_fn fn, void* arg); +void topo_fire_nodeinfo_cbk(struct UTUN_INSTANCE* instance, struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* node); + // ETCP ID для пакетов топологии #define ETCP_ID_TOPO_ENTRY 0x01 diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index bc8b7ce8..012e2940 100644 --- a/src/routing_layer/topo_node.c +++ b/src/routing_layer/topo_node.c @@ -53,6 +53,19 @@ 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; } } +/* извлекает 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) { + if (!nq || !nq->paths || !nq->paths->head) { *out_count = 0; if (out_rtt) *out_rtt = 0; return NULL; } + struct TOPO_NODEPATH* best = NULL; uint8_t min_hops = 255; + struct ll_entry* e = nq->paths->head; + while (e) { struct TOPO_NODEPATH* p = (struct TOPO_NODEPATH*)e; if (p->conn && p->conn->links_up && p->hop_count < min_hops) { best = p; min_hops = p->hop_count; } e = e->next; } + if (!best) best = (struct TOPO_NODEPATH*)nq->paths->head; + *out_count = best->hop_count; + if (out_rtt) *out_rtt = best->cumulative_rtt; + return (uint64_t*)((uint8_t*)best + sizeof(struct TOPO_NODEPATH)); +} + void topo_node_free_raw(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni) { if (!ni) return; if (groups) { @@ -87,7 +100,6 @@ void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_ u_free(nq->subnets); nq->subnets = NULL; } - u_free(nq->hop_list); nq->hop_list = NULL; nq->hop_count = 0; nq->handle = NULL; } @@ -169,7 +181,8 @@ 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.hop_count = nq->hop_count; + { uint8_t hc; uint64_t* hl = topo_node_best_hop_list(nq, &hc, NULL); + msg.hop_count = hc; } msg.cumulative_rtt = cumulative_rtt; size_t total = TOPOMSG_NODE_HDR_SIZE + (size_t)topo_node_dyn_size(&msg); @@ -208,10 +221,8 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_ }} } } - if (msg.hop_count && nq->hop_list) { - size_t sz = msg.hop_count * 8; - memcpy(dp, nq->hop_list, sz); dp += sz; - } + { uint8_t hc; uint64_t* hl = topo_node_best_hop_list(nq, &hc, NULL); + if (hc && hl) { size_t sz = hc * 8; memcpy(dp, hl, sz); dp += sz; } } return (int)(dp - out); } @@ -364,11 +375,22 @@ void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uin (unsigned long long)node_id, (unsigned)rtt, (unsigned long long)now, (unsigned long long)g->group_id); if (groups->instance && groups->instance->topo_sqlite_db) topo_node_sqlite_update_rtt(groups->instance->topo_sqlite_db, node_id, rtt); + topo_fire_nodeinfo_cbk(groups->instance, g, nq); } ge = ge->next; } } +uint16_t node_best_rtt(struct TOPO_GROUP_NODE* nq) { + if (!nq) return 0xFFFF; + struct TOPO_CONNECTIVITY* c = &nq->connectivity; + uint16_t best = 0x3FFF; uint8_t type = 0; + if (c->interface_status == PROBE_RESULT_REACHABLE && c->interface_min_rtt < best) { best = c->interface_min_rtt; type = 0; } + if (c->nat_status == PROBE_RESULT_REACHABLE && c->nat_min_rtt < best) { best = c->nat_min_rtt; type = 1; } + if (c->real_status == PROBE_RESULT_REACHABLE && c->real_min_rtt < best) { best = c->real_min_rtt; type = 2; } + return ((uint16_t)type << 14) | best; +} + // ===== dump / format ===== static int is_node_connectivity_active(struct TOPO_CONNECTIVITY* c) { @@ -658,7 +680,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->last_ver = ni->ver; e_sock = instance->etcp_sockets; diff --git a/src/routing_layer/topo_node.h b/src/routing_layer/topo_node.h index 5491e760..f8b6bc5d 100644 --- a/src/routing_layer/topo_node.h +++ b/src/routing_layer/topo_node.h @@ -66,6 +66,10 @@ struct UTUN_INSTANCE; #define CONN_TYPE_INDIRECT 3 #define CONN_MGR_MAX_INTERMEDIARIES 3 +#define NCONN_DIRECT (1 << 0) // прямое ETCP подключение +#define NCONN_INDIRECT (1 << 1) // через посредника +#define NCONN_BGP (1 << 2) // узел виден в топологии (BGP) + #define TOPO_FLAG_SEND_SUBNETS 0x01 #define TOPO_GROUP_UTUN 0x8000000000000000ULL #define TOPO_GROUP_TYPE_UTUN 1 @@ -151,15 +155,14 @@ struct TOPO_GROUP_NODE { 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) + uint8_t conn_presence; // NODE_CONN_* — какие подключения есть в принципе + uint8_t conn_up; // NODE_CONN_* — какие из них подняты struct CONN_MGR_HANDLE* handle; // активный хендл соединения к узлу }; @@ -179,6 +182,7 @@ int topo_node_serialize(struct TOPO_NODE* ni, struct TOPO_GROUP_NODE* nq, uint8_ 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, uint64_t** out_hop_list, uint8_t* out_hop_count, uint16_t* out_cumulative_rtt); +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); @@ -187,6 +191,9 @@ void topo_node_dump_all(struct TOPO_GROUP* group); int topo_node_format_all(struct TOPO_GROUP* group, char* buf, size_t buf_size); int topo_node_ping_request_cbk(struct TOPO_GROUPS* groups, uint64_t node_id); void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uint16_t rtt); +uint16_t node_best_rtt(struct TOPO_GROUP_NODE* nq); +#define NODE_RTT_TYPE(v) ((v) >> 14) +#define NODE_RTT_VALUE(v) ((v) & 0x3FFF) static inline const struct TOPO_SOCKMETA4* topo_v4_sock_meta(const struct TOPO_NODE* ni) { return ni->v4_sock_meta; } static inline const struct TOPO_ADDR4* topo_v4_addrs(const struct TOPO_NODE* ni) { return ni->v4_addrs; } diff --git a/src/routing_layer/topo_node_sqlite.c b/src/routing_layer/topo_node_sqlite.c index 264ee97c..1cb31153 100644 --- a/src/routing_layer/topo_node_sqlite.c +++ b/src/routing_layer/topo_node_sqlite.c @@ -58,6 +58,7 @@ int topo_node_sqlite_init(sqlite3* db) { ");" "CREATE INDEX IF NOT EXISTS idx_na_node ON node_addresses(node_id);" "CREATE UNIQUE INDEX IF NOT EXISTS idx_na_unique ON node_addresses(node_id, family, socket_id, addr_type);" + "CREATE INDEX IF NOT EXISTS idx_na_addr_lookup ON node_addresses(family, address, port, socket_id);" "CREATE TABLE IF NOT EXISTS channels (" " channel_id TEXT PRIMARY KEY," // ID канала @@ -178,6 +179,51 @@ int topo_node_sqlite_node_put(sqlite3* db, struct TOPO_GROUPS* groups, struct TO } } + /* Удалить адреса других узлов, конфликтующие с адресами текущего */ + { + sqlite3_stmt* del_cfl = NULL; + if (sqlite3_prepare_v2(db, + "DELETE FROM node_addresses WHERE node_id != ?1 AND family = ?2 AND address = ?3 AND port = ?4 AND socket_id = ?5", + -1, &del_cfl, NULL) == SQLITE_OK) { + struct TOPO_ADDR4* a4 = ni->v4_addrs; + while (a4) { + sqlite3_bind_int64(del_cfl, 1, (sqlite3_int64)ni->node_id); + sqlite3_bind_int(del_cfl, 2, 4); + sqlite3_bind_blob(del_cfl, 3, a4->addr, 4, SQLITE_STATIC); + sqlite3_bind_int(del_cfl, 4, a4->port); + sqlite3_bind_int(del_cfl, 5, (int)a4->socket_id); + sqlite3_step(del_cfl); + if (sqlite3_changes(db) > 0) + DEBUG_INFO(DEBUG_CATEGORY_BGP, "node_put: removed stale v4 addr %d.%d.%d.%d:%u sock=%d other_node nid=%016llx", + a4->addr[0], a4->addr[1], a4->addr[2], a4->addr[3], a4->port, (int)a4->socket_id, + (unsigned long long)ni->node_id); + sqlite3_reset(del_cfl); + a4 = a4->next; + } + struct TOPO_ADDR6* a6 = ni->v6_addrs; + while (a6) { + if (a6->addr[0] == 0xFE && (a6->addr[1] & 0xC0) == 0x80) { a6 = a6->next; continue; } + sqlite3_bind_int64(del_cfl, 1, (sqlite3_int64)ni->node_id); + sqlite3_bind_int(del_cfl, 2, 6); + sqlite3_bind_blob(del_cfl, 3, a6->addr, 16, SQLITE_STATIC); + sqlite3_bind_int(del_cfl, 4, a6->port); + sqlite3_bind_int(del_cfl, 5, (int)a6->socket_id); + sqlite3_step(del_cfl); + sqlite3_reset(del_cfl); + a6 = a6->next; + } + sqlite3_finalize(del_cfl); + } + } + + addr_stmt_beg: + { + int v4cnt = topo_list_count((struct _topo_head*)ni->v4_addrs); + int v6cnt = topo_list_count((struct _topo_head*)ni->v6_addrs); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "node_put: nid=%016llx INSERT v4=%d v6=%d", + (unsigned long long)ni->node_id, v4cnt, v6cnt); + } + sqlite3_stmt* addr_stmt = NULL; if (sqlite3_prepare_v2(db, "INSERT INTO node_addresses(node_id, family, protocol, address, port, rtt, addr_type, socket_id)" @@ -240,13 +286,25 @@ int topo_node_sqlite_node_put(sqlite3* db, struct TOPO_GROUPS* groups, struct TO sqlite3_bind_null(addr_stmt, 6); sqlite3_bind_int(addr_stmt, 7, at); sqlite3_bind_int(addr_stmt, 8, (int)a6->socket_id); - sqlite3_step(addr_stmt); sqlite3_reset(addr_stmt); + { int rc = sqlite3_step(addr_stmt); + if (rc != SQLITE_DONE && rc != SQLITE_OK) + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "node_put: INSERT v6 FAILED nid=%016llx rc=%d err=%s", + (unsigned long long)ni->node_id, rc, sqlite3_errmsg(db)); } + sqlite3_reset(addr_stmt); a6 = a6->next; } sqlite3_finalize(addr_stmt); } sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + { int total = 0; + sqlite3_stmt* cnt_st; + if (sqlite3_prepare_v2(db, "SELECT count(*) FROM node_addresses WHERE node_id=?", -1, &cnt_st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(cnt_st, 1, (sqlite3_int64)ni->node_id); + if (sqlite3_step(cnt_st) == SQLITE_ROW) total = sqlite3_column_int(cnt_st, 0); + sqlite3_finalize(cnt_st); } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "node_put: nid=%016llx COMMITED total=%d", + (unsigned long long)ni->node_id, total); } return 0; } diff --git a/src/routing_layer/topo_recovery.c b/src/routing_layer/topo_recovery.c index a9c680f8..a44d0e8b 100644 --- a/src/routing_layer/topo_recovery.c +++ b/src/routing_layer/topo_recovery.c @@ -72,16 +72,6 @@ static uint16_t node_get_min_rtt(const struct TOPO_GROUP_NODE* nq) { return rtt; } -/* Вычисляет next_hop относительно failed_peer из hop_list. - * hop_list = [...узлы пути...], последний элемент = failed_peer (форвардил NODEINFO нам). - * next_hop = элемент перед failed_peer (hop_list[hop_count-2]). - * Если failed_peer не найден в конце — fallback: node_id сам себе next_hop. */ -static uint64_t find_next_hop(const struct TOPO_GROUP_NODE* nq, uint64_t failed_peer) { - if (nq->hop_count >= 2 && nq->hop_list && nq->hop_list[nq->hop_count - 1] == failed_peer) - return nq->hop_list[nq->hop_count - 2]; - return nq->node_id; -} - /* Ищет незапущенный (started==0) контекст по next_hop_id */ static struct TOPO_RECOVERY_CTX* topo_recovery_find_by_next_hop(struct TOPO_GROUP* group, uint64_t next_hop_id) { struct TOPO_RECOVERY_CTX* ctx = group->recovery_list; @@ -183,12 +173,11 @@ static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx) { } } -void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t failed_peer) { +void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t next_hop) { if (!group || !nq) return; struct TOPO_NODE* ni = topo_node_registry_find(group->instance->topo_groups, nq->node_id); if (!ni) return; uint64_t node_id = nq->node_id; - uint64_t next_hop = find_next_hop(nq, failed_peer); struct TOPO_RECOVERY_CTX* ctx = topo_recovery_find_by_next_hop(group, next_hop); if (!ctx) { diff --git a/src/routing_layer/topo_recovery.h b/src/routing_layer/topo_recovery.h index c3baee5a..9cf50e83 100644 --- a/src/routing_layer/topo_recovery.h +++ b/src/routing_layer/topo_recovery.h @@ -59,7 +59,7 @@ struct TOPO_RECOVERY_CTX { * Вызывается из topo_group_remove_conn ДО topo_node_free_lists (нужен hop_list). * @param failed_peer node_id отвалившегося пира (conn->peer_node_id) */ -void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t failed_peer); +void topo_recovery_add_node(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, uint64_t next_hop); /** * Запускает перебор узлов для всех pending-контекстов. diff --git a/src/transport_layer/node_conn_direct.c b/src/transport_layer/node_conn_direct.c index 82d03dd9..1efa915c 100644 --- a/src/transport_layer/node_conn_direct.c +++ b/src/transport_layer/node_conn_direct.c @@ -127,7 +127,28 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni, struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); struct sockaddr_storage sa; memcpy(&sa, &sin, sizeof(sin)); - if (etcp_link_new(conn, socks[rr++ % sock_count], &sa, 0)) link_count++; + struct ETCP_SOCKET* use_sock = socks[rr++ % sock_count]; + { + struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa); + if (stale && stale->etcp != conn + && memcmp(conn->crypto_ctx.peer_public_key, stale->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) + { + if (!stale->etcp->links_up || !stale->etcp->initialized) { + struct ETCP_CONN* old = stale->etcp; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[ncd] evict stale link node=0x%016llx(old pk=%016llx) for new node=0x%016llx(pk=%016llx)", + (unsigned long long)old->peer_node_id, *(const uint64_t*)old->crypto_ctx.peer_public_key, + (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); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "[ncd] addr conflict: %d.%d.%d.%d:%d already used by live conn 0x%016llx for node 0x%016llx", + a->addr[0], a->addr[1], a->addr[2], a->addr[3], (int)a->port, + (unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id); + continue; + } + } + } + if (etcp_link_new(conn, use_sock, &sa, 0)) link_count++; } } } @@ -155,6 +176,24 @@ static int ncd_create_links(struct ncd_entry* entry, struct TOPO_NODE* ni, struct ETCP_SOCKET* use_sock = socks[rr++ % sock_count]; if ((a->addr[0] == 0xfe && (a->addr[1] & 0xc0) == 0x80)) sin6.sin6_scope_id = use_sock->netif_index; struct sockaddr_storage sa; memcpy(&sa, &sin6, sizeof(sin6)); + { + struct ETCP_LINK* stale = etcp_link_find_by_addr(use_sock, &sa); + if (stale && stale->etcp != conn + && memcmp(conn->crypto_ctx.peer_public_key, stale->etcp->crypto_ctx.peer_public_key, SC_PUBKEY_SIZE)) + { + if (!stale->etcp->links_up || !stale->etcp->initialized) { + struct ETCP_CONN* old = stale->etcp; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[ncd] evict stale v6 link node=0x%016llx(old) -> 0x%016llx(new)", + (unsigned long long)old->peer_node_id, (unsigned long long)conn->peer_node_id); + etcp_link_close(stale); + etcp_cbk_fire(old, ETCP_CBK_EVENT_NODE_CHANGED); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "[ncd] addr conflict: v6 port=%d already used by live conn 0x%016llx for node 0x%016llx", + (int)a->port, (unsigned long long)stale->etcp->peer_node_id, (unsigned long long)conn->peer_node_id); + continue; + } + } + } if (etcp_link_new(conn, use_sock, &sa, 0)) link_count++; } } @@ -183,8 +222,9 @@ static void ncd_event_dispatch(struct ncd_entry* entry, enum ncd_event event) { entry->connect_timer = NULL; break; } - DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] event=%d node=0x%016llx handles=%d", - (int)event, (unsigned long long)entry->node_id, entry->handle_count); + DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] event=%s node=0x%016llx handles=%d", + event == NCD_EVENT_UP ? "UP" : event == NCD_EVENT_DOWN ? "DOWN" : "TIMEOUT", + (unsigned long long)entry->node_id, entry->handle_count); struct NODE_CONN_DIRECT* h = entry->handles; while (h) { struct NODE_CONN_DIRECT* next = h->next; if (h->cb) h->cb(h, event, h->cb_arg); h = next; } } @@ -220,8 +260,6 @@ static void ncd_deferred_close_conn(void* arg) { static void ncd_down_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)event; struct ncd_entry* entry = (struct ncd_entry*)arg; if (!entry) return; - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "[ncd-debug] ncd_down_cb conn=%p entry=%p fin_wait=%d handle_count=%d node=0x%016llx", - (void*)conn, entry, conn->fin_wait, entry->handle_count, (unsigned long long)entry->node_id); if (conn->fin_wait && entry->handle_count <= 0) { DEBUG_INFO(DEBUG_CATEGORY_NCD, "[ncd] DOWN during fin_wait, cleaning up node=0x%016llx", (unsigned long long)entry->node_id); conn->fin_wait = 0; conn->fin_wait_clear_cb = NULL; conn->fin_wait_clear_arg = NULL; diff --git a/src/transport_layer/stcp_link.c b/src/transport_layer/stcp_link.c index 3e06519c..d17f2841 100644 --- a/src/transport_layer/stcp_link.c +++ b/src/transport_layer/stcp_link.c @@ -173,7 +173,8 @@ struct stcp_link *stcp_link_connect(struct stcp_link_config *cfg) { link->etcp_conn.instance = cfg->inst; link->etcp_conn.transport_link = (void *)link; - DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "connecting to %s:%u", addr_str, port); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "stcp_link: connecting to %s:%u pubkey=%016llx", + addr_str, port, (unsigned long long)*(const uint64_t*)link->peer_pubkey); return link; } diff --git a/src/utun_instance.c b/src/utun_instance.c index 9a09dbd5..5515b464 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -618,6 +618,13 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Free the instance memory DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Freeing instance memory"); + + // Cleanup nodeinfo callback chain + while (instance->nodeinfo_cbks) { + struct nodeinfo_cbk_entry* r = instance->nodeinfo_cbks; + instance->nodeinfo_cbks = r->next; u_free(r); + } + u_free(instance); DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] Instance destroyed completely"); } diff --git a/src/utun_instance.h b/src/utun_instance.h index 6e25322b..59cb83b0 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -28,6 +28,7 @@ extern "C" { #include "etcp_api.h" #include "config_parser.h" +#include "topo_group.h" // Forward declarations struct utun_config; @@ -126,6 +127,7 @@ struct UTUN_INSTANCE { struct etcp_inst_cbk_entry* new_conn_cbks; struct etcp_status_cbk_entry* conn_status_cbks; // instance-level: NEW/UP/DOWN/DELETE struct etcp_socket_cbk_entry* socket_cbks; // instance-level: socket ADDR/STATUS changes + struct nodeinfo_cbk_entry* nodeinfo_cbks; // глобальная подписка на изменения nodeinfo void* test_user_ptr; // Generic user pointer (used by tests) struct memory_pool* data_pool;// для входных-выходных данных пакета diff --git a/tests/test_route_lib.c b/tests/test_route_lib.c index 78811ee4..812922ce 100644 --- a/tests/test_route_lib.c +++ b/tests/test_route_lib.c @@ -50,7 +50,6 @@ static struct TOPO_GROUP_NODE* create_test_node(uint64_t node_id, uint8_t ver, u ni->node_id = htobe64(node_id); ni->ver = ver; nq->node_id = ni->node_id; - nq->hop_count = 0; nq->paths = queue_new(NULL, 16, 0, 0, "test_paths"); struct TOPO_NODESUBNETS* r = u_calloc(1, sizeof(struct TOPO_NODESUBNETS)); if (r) { @@ -108,7 +107,7 @@ static void test_change_cb(struct ROUTE_TABLE* table, static void test_table_create_destroy(void) { TEST("route_table_create / destroy"); - struct ROUTE_TABLE *t = route_table_create(); + struct ROUTE_TABLE *t = route_table_create(0); ASSERT_PTR(t, "create failed"); ASSERT_EQ(t->count, 0, "empty table"); route_table_destroy(t); @@ -117,7 +116,7 @@ static void test_table_create_destroy(void) { static void test_nodeinfo_insert(void) { TEST("TOPO_GROUP_NODE insert into routing table"); - struct ROUTE_TABLE *t = route_table_create(); + struct ROUTE_TABLE *t = route_table_create(0); struct TOPO_GROUP_NODE* nq = create_test_node(0x12345678ULL, 1, 0x0a000000, 8); bool inserted = route_insert(t, nq); @@ -131,7 +130,7 @@ static void test_nodeinfo_insert(void) { static void test_versioning(void) { TEST("NODEINFO multiple nodes"); - struct ROUTE_TABLE *t = route_table_create(); + struct ROUTE_TABLE *t = route_table_create(0); struct TOPO_GROUP_NODE* nq1 = create_test_node(0x1111ULL, 5, 0x0a000000, 8); struct TOPO_GROUP_NODE* nq2 = create_test_node(0x2222ULL, 1, 0xac100000, 12); @@ -150,7 +149,7 @@ static void test_versioning(void) { static void test_performance_1000_nodes(void) { TEST("performance with 1000 nodes"); - struct ROUTE_TABLE *t = route_table_create(); + struct ROUTE_TABLE *t = route_table_create(0); clock_t start = clock(); for (int i = 0; i < 1000; i++) { @@ -172,7 +171,7 @@ static void test_performance_1000_nodes(void) { static void test_destroy_refcount(void) { TEST("destroy with multiple paths (ref_count safety)"); - struct ROUTE_TABLE *t = route_table_create(); + struct ROUTE_TABLE *t = route_table_create(0); /* TOPO_GROUP_NODE memory managed externally, table only holds pointers */ route_table_destroy(t); PASS(); diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt index 2d19d465..588b3a3a 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt @@ -3,6 +3,47 @@ package com.utun.chat.data import android.util.Log import org.json.JSONArray import org.json.JSONObject +import java.nio.ByteBuffer +import java.nio.ByteOrder + +const val EVT_NODEINFO_UPDATED = 16 + +data class NodeStatus( + val nodeId: Long = 0, + val directPresence: Boolean = false, + val directUp: Boolean = false, + val indirectPresence: Boolean = false, + val indirectUp: Boolean = false, + val bgpPresence: Boolean = false, + val bgpUp: Boolean = false, + val bestRtt: Int = 0, + val rttType: Int = 0 // 0=interface, 1=nat, 2=real +) { + val online: Boolean get() = (directPresence && directUp) || (indirectPresence && indirectUp) + val hasPresence: Boolean get() = directPresence || indirectPresence || bgpPresence +} + +fun parseNodeStatus(data: ByteArray): NodeStatus? { + if (data.size < 12) return null + val buf = ByteBuffer.wrap(data, 0, 12).order(ByteOrder.LITTLE_ENDIAN) + val nodeId = buf.getLong() + val presence = buf.get().toInt() and 0xFF + val up = buf.get().toInt() and 0xFF + val rttPacked = buf.getShort().toInt() and 0xFFFF + val bestRtt = rttPacked and 0x3FFF + val rttType = (rttPacked shr 14) and 3 + return NodeStatus( + nodeId = nodeId, + directPresence = (presence and 1) != 0, + directUp = (up and 1) != 0, + indirectPresence = (presence and 2) != 0, + indirectUp = (up and 2) != 0, + bgpPresence = (presence and 4) != 0, + bgpUp = (up and 4) != 0, + bestRtt = bestRtt, + rttType = rttType + ) +} data class Channel(val id: String, val name: String, val lastMsgAt: Long = 0, val lastMsg: String = "", val peersOnline: Int = 0) @@ -15,11 +56,18 @@ data class Message(val id: Long, val author: String, val text: String, val ts: L data class ChatMember(val nodeId: Long, val name: String, val online: Boolean = false, val isAdmin: Boolean = false, val isSuper: Boolean = false, val isStorage: Boolean = false, val nodeType: Int = 0, - val rtt: Int = 0, val isSelf: Boolean = false) + val rtt: Int = 0, val rttType: Int = 0, val isSelf: Boolean = false, + val directUp: Boolean = false, val indirectUp: Boolean = false, + val directPresence: Boolean = false, val indirectPresence: Boolean = false, + val bgpPresence: Boolean = false) data class MemberDetail(val nodeId: Long, val name: String = "", val online: Boolean = false, val lastSeen: Long = 0, - val created: Long = 0, val addresses: List = emptyList()) + val created: Long = 0, val addresses: List = emptyList(), + val bestRtt: Int = 0, val rttType: Int = 0, + val directPresence: Boolean = false, val directUp: Boolean = false, + val indirectPresence: Boolean = false, val indirectUp: Boolean = false, + val bgpPresence: Boolean = false) data class MemberAddr(val family: Int = 0, val protocol: String = "", val address: String = "", val port: Int = 0, val socketId: Int = 0, val rtt: Int = 0) diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/MemberListScreen.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/MemberListScreen.kt index 23aba9cc..21d5e7e5 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/MemberListScreen.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/MemberListScreen.kt @@ -121,7 +121,17 @@ private fun MemberListItem( onClick: () -> Unit, onEdit: () -> Unit ) { - val onlineColor = if (member.online) Color(0xFF4CAF50) else Color(0xFF9E9E9E) + val onlineColor = when { + member.directUp || member.indirectUp -> Color(0xFF4CAF50) // green: up + member.directPresence || member.indirectPresence || member.bgpPresence -> Color(0xFFFFA726) // orange: presence but down + else -> Color(0xFF9E9E9E) // gray: offline + } + val connLabel = when { + member.directUp -> "direct" + member.indirectUp -> "indirect" + member.directPresence || member.indirectPresence || member.bgpPresence -> "offline" + else -> "" + } val bgColor = if (isSelected) MaterialTheme.colorScheme.primaryContainer else Color.Transparent Surface( @@ -152,27 +162,22 @@ private fun MemberListItem( overflow = TextOverflow.Ellipsis, color = if (member.isSelf) Color(0xFF4CAF50) else Color.Unspecified ) - val roles = buildList { - if (member.isAdmin) add("Admin") - if (member.isSuper) add("Super") - if (member.isStorage) add("Storage") - } - if (roles.isNotEmpty() || member.rtt > 0) { + val parts = mutableListOf() + if (member.isAdmin) parts.add("Admin") + if (member.isSuper) parts.add("Super") + if (member.isStorage) parts.add("Storage") + if (member.directUp) parts.add("direct↑") + else if (member.indirectUp) parts.add("indirect↑") + else if (member.directPresence || member.indirectPresence || member.bgpPresence) parts.add("down") + if (member.rtt > 0) parts.add("${member.rtt}ms") + if (parts.isNotEmpty()) { Row(verticalAlignment = Alignment.CenterVertically) { - if (roles.isNotEmpty()) - Text( - text = roles.joinToString(", "), - fontSize = 11.sp, - color = MaterialTheme.colorScheme.primary - ) - if (member.rtt > 0) { - if (roles.isNotEmpty()) Text(" ", fontSize = 11.sp) - Text( - text = "${member.rtt / 10}ms", - fontSize = 11.sp, - color = Color.Gray - ) - } + Text( + text = parts.joinToString(" "), + fontSize = 11.sp, + color = if (member.directUp || member.indirectUp) + Color(0xFF4CAF50) else Color.Gray + ) } } } @@ -266,6 +271,38 @@ private fun MemberDetailCard( } } } + val rttTypeName = when (detail.rttType) { 0 -> "interface" 1 -> "nat" 2 -> "real" else -> "?" } + Spacer(Modifier.height(6.dp)) + Text("Connection:", fontSize = 13.sp, fontWeight = FontWeight.Medium) + if (detail.directPresence) { + Row(verticalAlignment = Alignment.CenterVertically) { + Box(Modifier.size(8.dp).clip(CircleShape) + .background(if (detail.directUp) Color(0xFF4CAF50) else Color(0xFFE53935))) + Spacer(Modifier.width(6.dp)) + Text("Direct: ${if (detail.directUp) "UP" else "DOWN"}" + + (if (detail.directUp && detail.bestRtt > 0) " ${detail.bestRtt}ms ($rttTypeName)" else ""), + fontSize = 12.sp, color = Color.DarkGray) + } + } + if (detail.indirectPresence) { + Row(verticalAlignment = Alignment.CenterVertically) { + Box(Modifier.size(8.dp).clip(CircleShape) + .background(if (detail.indirectUp) Color(0xFF4CAF50) else Color(0xFFE53935))) + Spacer(Modifier.width(6.dp)) + Text("Indirect: ${if (detail.indirectUp) "UP" else "DOWN"}", + fontSize = 12.sp, color = Color.DarkGray) + } + } + if (detail.bgpPresence) { + Row(verticalAlignment = Alignment.CenterVertically) { + Box(Modifier.size(8.dp).clip(CircleShape).background(Color(0xFF4CAF50))) + Spacer(Modifier.width(6.dp)) + Text("BGP: visible", fontSize = 12.sp, color = Color.DarkGray) + } + } + if (!detail.directPresence && !detail.indirectPresence && !detail.bgpPresence) { + Text("offline", fontSize = 12.sp, color = Color.Gray) + } } } } diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt index 79020071..1fd5d2eb 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt @@ -57,6 +57,9 @@ class ChatViewModel : ViewModel() { private val _isStuck = MutableStateFlow(false) val isStuck: StateFlow = _isStuck + private val _nodeStatuses = MutableStateFlow>(emptyMap()) + val nodeStatuses: StateFlow> = _nodeStatuses + val audioRecorder = AudioRecorderManager() private val _isRecording = MutableStateFlow(false) @@ -186,6 +189,13 @@ class ChatViewModel : ViewModel() { 4 -> if (repo != null) refreshChannels() /* CHANNEL_UPDATED */ 12 -> LogManager.addLog("INFO", "VM", "service started") 13 -> LogManager.addLog("INFO", "VM", "service stopped") + 16 -> { /* NODEINFO_UPDATED: [node_id:8][presence:1][up:1][best_rtt:2] */ + if (data == null || data.size < 12) return + val ns = parseNodeStatus(data) ?: return + _nodeStatuses.value = _nodeStatuses.value + (ns.nodeId to ns) + val cur = _currentChannel.value + if (cur != null) refreshMembers(cur.id) + } } } @@ -228,15 +238,33 @@ class ChatViewModel : ViewModel() { fun refreshMembers(chId: String) { val r = repo ?: return + val statuses = _nodeStatuses.value viewModelScope.launch { - _members.value = withContext(Dispatchers.IO) { r.getMembers(chId) } + val raw = withContext(Dispatchers.IO) { r.getMembers(chId) } + _members.value = raw.map { m -> + val ns = statuses[m.nodeId] + if (ns != null) m.copy( + online = ns.online, rtt = ns.bestRtt, rttType = ns.rttType, + directUp = ns.directUp, indirectUp = ns.indirectUp, + directPresence = ns.directPresence, indirectPresence = ns.indirectPresence, + bgpPresence = ns.bgpPresence + ) else m + } } } fun showMemberDetail(channelId: String, nodeId: Long) { val r = repo ?: return + val ns = _nodeStatuses.value[nodeId] viewModelScope.launch { - _selectedMember.value = withContext(Dispatchers.IO) { r.getMemberDetail(channelId, nodeId) } + val d = withContext(Dispatchers.IO) { r.getMemberDetail(channelId, nodeId) } + if (d != null && ns != null) _selectedMember.value = d.copy( + online = ns.online, bestRtt = ns.bestRtt, rttType = ns.rttType, + directPresence = ns.directPresence, directUp = ns.directUp, + indirectPresence = ns.indirectPresence, indirectUp = ns.indirectUp, + bgpPresence = ns.bgpPresence + ) + else _selectedMember.value = d } } diff --git a/tools/chatgui-android/headless/headless_control.c b/tools/chatgui-android/headless/headless_control.c index 1e32fa8b..9633cc79 100644 --- a/tools/chatgui-android/headless/headless_control.c +++ b/tools/chatgui-android/headless/headless_control.c @@ -5,6 +5,8 @@ #include "../libutun_lite/utun_config_api.h" #include "../libutun_lite/invite_link_c.h" #include "../../../src/chat/chat_setting.h" +#include "../../../src/chat/chat_sync.h" +#include "../../../lib/debug_config.h" #include #include @@ -291,16 +293,35 @@ static int handle_join(int fd, int id, const char* json) { send_response(fd, id, NULL, err); return 0; } + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "join: ch=%llu node=0x%016llx pubkey=%02x%02x%02x%02x... addrs=%d", + (unsigned long long)d.channelId, (unsigned long long)d.nodeId, + d.pubkey[0], d.pubkey[1], d.pubkey[2], d.pubkey[3], (int)d.addrCount); + + for (int i = 0; i < (int)d.addrCount; i++) { + struct InviteAddrC* a = &d.addrs[i]; + char ip[64]; snprintf(ip, sizeof(ip), a->family == 6 ? "v6" : "%d.%d.%d.%d", + a->address[0], a->address[1], a->address[2], a->address[3]); + const char* proto = a->proto == 1 ? "UDP" : a->proto == 2 ? "TCP" : a->proto == 3 ? "UDP+TCP" : "?"; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, " addr[%d]: %s:%d %s sock=%d", i, ip, (int)a->port, proto, (int)a->socketId); + } + + uint8_t addrs_buf[512]; + int addrs_len = invite_serialize_addrs(&d, addrs_buf, sizeof(addrs_buf)); + if (addrs_len <= 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "join: invite_serialize_addrs failed"); + send_response(fd, id, NULL, "internal error: serialize failed"); + return 0; + } + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "join: calling chat_sync_connect_from_invite..."); + chat_sync_connect_from_invite(d.channelId, d.nodeId, d.pubkey, addrs_buf, d.addrCount, addrs_len); + char resp[512]; snprintf(resp, sizeof(resp), - "{\"channel_id\":%llu,\"node_id\":\"0x%llx\",\"addrs\":%d}", + "{\"channel_id\":%llu,\"node_id\":\"0x%llx\",\"addrs\":%d,\"connecting\":true}", (unsigned long long)d.channelId, (unsigned long long)d.nodeId, (int)d.addrCount); send_response(fd, id, resp, NULL); - - /* TODO: call chat_sync_connect_from_invite() once integrated */ - printf("control: join ch=%llu node=0x%llx pubkey=%016llx... addrs=%d\n", - (unsigned long long)d.channelId, (unsigned long long)d.nodeId, - (unsigned long long)*(const uint64_t*)d.pubkey, (int)d.addrCount); return 0; } diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.c b/tools/chatgui-android/jni_bridge/android_jni_bridge.c index b87a43d6..6cdafc52 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.c +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.c @@ -213,15 +213,16 @@ void utun_bridge_connect_node(const char* address, int port, const char* pubkey_ void utun_bridge_join_channel(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, - const uint8_t* addrs_data, int addr_count) { - bridge_log(BLEV_INFO, "join ch=%llu node=0x%016llx addrs=%d", - (unsigned long long)channel_id, (unsigned long long)node_id, addr_count); + const uint8_t* addrs_data, int addr_count, + int addrs_data_len) { + bridge_log(BLEV_INFO, "join ch=%llu node=0x%016llx addrs=%d len=%d", + (unsigned long long)channel_id, (unsigned long long)node_id, addr_count, addrs_data_len); if (!pubkey_bin || !addrs_data || addr_count <= 0) { bridge_log(BLEV_ERROR, "join invalid args"); utun_bridge_event(6, "{\"node_id\":0,\"result\":-7}"); /* CHAT_EVT_CONNECT_RESULT */ return; } - chat_sync_connect_from_invite(channel_id, node_id, pubkey_bin, addrs_data, addr_count); + chat_sync_connect_from_invite(channel_id, node_id, pubkey_bin, addrs_data, addr_count, addrs_data_len); } void utun_bridge_connect_channel(const char* channel_id) { @@ -1125,7 +1126,7 @@ JNIEXPORT void JNICALL Java_com_utun_chat_data_NativeLib_nativeJoinChannel( uint8_t* ad = (uint8_t*)(*env)->GetByteArrayElements(env, addrsData, NULL); bridge_log(BLEV_DEBUG, "nativeJoinChannel ch=%lld node=0x%llx pk=%02x%02x... addrs_sz=%d", (long long)channelId, (unsigned long long)nodeId, pk[0], pk[1], (int)ad_len); - utun_bridge_join_channel((uint64_t)channelId, (uint64_t)nodeId, pk, ad, (int)addrCount); + utun_bridge_join_channel((uint64_t)channelId, (uint64_t)nodeId, pk, ad, (int)addrCount, (int)ad_len); (*env)->ReleaseByteArrayElements(env, addrsData, (jbyte*)ad, JNI_ABORT); (*env)->ReleaseByteArrayElements(env, pubkey, (jbyte*)pk, JNI_ABORT); } diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.h b/tools/chatgui-android/jni_bridge/android_jni_bridge.h index 6fc7a067..b58b0031 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.h +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.h @@ -54,7 +54,8 @@ void utun_bridge_connect_node(const char* address, int port, const char* pubkey_ addr_data format: family(1) + socketId(1) + address(4|16) + port(2 BE) per addr */ void utun_bridge_join_channel(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, - const uint8_t* addrs_data, int addr_count); + const uint8_t* addrs_data, int addr_count, + int addrs_data_len); /* Connect to all peers of a channel (auto-connect on select) */ void utun_bridge_connect_channel(const char* channel_id); diff --git a/tools/chatgui-android/libutun_lite/instance_lite.c b/tools/chatgui-android/libutun_lite/instance_lite.c index 52521f9f..db11a5cb 100644 --- a/tools/chatgui-android/libutun_lite/instance_lite.c +++ b/tools/chatgui-android/libutun_lite/instance_lite.c @@ -134,6 +134,18 @@ static void chat_event_forward(int type, const uint8_t* data, int len) { if (g_event_handler) g_event_handler(type, data, len); } +static void nodeinfo_event_cb(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, void* arg) { + (void)group; (void)arg; + if (!nq || !g_event_handler) return; + uint16_t best_rtt = node_best_rtt(nq); + uint8_t data[12]; + memcpy(data, &nq->node_id, 8); + data[8] = nq->conn_presence; + data[9] = nq->conn_up; + memcpy(data + 10, &best_rtt, 2); + g_event_handler(CHAT_EVT_NODEINFO_UPDATED, data, sizeof(data)); +} + /* ── Thread function ── */ static void* instance_thread(void* arg) { @@ -198,6 +210,8 @@ static void* instance_thread(void* arg) { /* Bind chat sync via etcp_router */ etcp_router_bind(g_inst, ETCP_RT_ID_CHAT_SYNC, NULL); + utun_add_nodeinfo_cbk(g_inst, nodeinfo_event_cb, NULL); + /* Set my_name from config */ if (g_inst->config->global.name[0]) { chat_core_update_my_name(g_inst->config->global.name); diff --git a/tools/chatgui/src/joindialog.cpp b/tools/chatgui/src/joindialog.cpp index 153d2292..75c3bce8 100644 --- a/tools/chatgui/src/joindialog.cpp +++ b/tools/chatgui/src/joindialog.cpp @@ -188,7 +188,7 @@ void JoinDialog::onConnectClicked() { chat_sync_connect_from_invite(d.channelId, nodeId, (const uint8_t*)d.pubkey.constData(), - (const uint8_t*)addrsBuf.constData(), addrCount); + (const uint8_t*)addrsBuf.constData(), addrCount, addrsBuf.size()); } void JoinDialog::onConnectResult(uint64_t nodeId, uint64_t channelId, int result) { diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h index 404af2d9..4b46a4b5 100644 --- a/tools/chatgui/transport/gui_bridge.h +++ b/tools/chatgui/transport/gui_bridge.h @@ -9,6 +9,8 @@ extern "C" { #include struct UASYNC; +struct TOPO_GROUP; +struct TOPO_GROUP_NODE; /* ── Типы уведомлений uasync→GUI (fire-and-forget) ── */ @@ -25,6 +27,7 @@ struct UASYNC; #define GUI_EVT_NODE_CHANGED 11 /* data: [node_id:8] */ #define GUI_EVT_ATTACHMENT_DOWNLOADED 14 /* data: [ch_id_len:1][ch_id:var][msg_id:8] */ #define GUI_EVT_DOWNLOAD_PROGRESS 15 /* data: [ch_id_len:1][ch_id:var][msg_id:8][blocks_done:4][num_blocks:4] */ +#define GUI_EVT_NODEINFO_UPDATE 16 /* data: node_id:8 conn_presence:1 conn_up:1 best_rtt_packed:2 */ /* ── API ── */ @@ -88,6 +91,13 @@ void gui_bridge_set_attachment_downloaded_cb(gui_attachment_downloaded_fn cb); typedef void (*gui_download_progress_fn)(const char* ch_id, int ch_id_len, int64_t msg_id, int blocks_done, int num_blocks); void gui_bridge_set_download_progress_cb(gui_download_progress_fn cb); +/* Callback для обновления статуса узла (nodeinfo) — data: 12B packed */ +typedef void (*gui_nodeinfo_update_fn)(const uint8_t* data, int len); +void gui_bridge_set_nodeinfo_update_cb(gui_nodeinfo_update_fn cb); + +/* uTun nodeinfo callback — registered via utun_add_nodeinfo_cbk, called from uasync thread */ +void gui_nodeinfo_cb_impl(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, void* arg); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp index 308d72fc..ed539a52 100644 --- a/tools/chatgui/transport/gui_bridge_impl.cpp +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -9,6 +9,8 @@ extern "C" { #include "../../../lib/u_async.h" #include "../../../lib/mem.h" #include "../../../lib/debug_config.h" +#include "../../../src/routing_layer/topo_node.h" +#include "../../../src/routing_layer/topo_group.h" } /* ── Внутренний объект-приёмник в GUI-потоке ── */ @@ -35,6 +37,7 @@ static gui_db_ready_fn g_db_ready_cb = nullptr; static gui_status_refresh_fn g_status_refresh_cb = nullptr; static gui_attachment_downloaded_fn g_attachment_downloaded_cb = nullptr; static gui_download_progress_fn g_download_progress_cb = nullptr; +static gui_nodeinfo_update_fn g_nodeinfo_update_cb = nullptr; static struct UASYNC* g_ua = nullptr; /* ── GuiBridgeReceiver implementation ── */ @@ -164,6 +167,9 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { } } break; + case GUI_EVT_NODEINFO_UPDATE: + if (g_nodeinfo_update_cb) g_nodeinfo_update_cb(d, dlen); + break; default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); break; @@ -254,6 +260,24 @@ void gui_bridge_set_download_progress_cb(gui_download_progress_fn cb) { g_download_progress_cb = cb; } +void gui_bridge_set_nodeinfo_update_cb(gui_nodeinfo_update_fn cb) { + g_nodeinfo_update_cb = cb; +} + } /* extern "C" */ +/* ── uTun nodeinfo callback (called from uasync thread) ── */ + +extern "C" void gui_nodeinfo_cb_impl(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* nq, void* arg) { + (void)arg; + if (!nq) return; + uint16_t best_rtt = node_best_rtt(nq); + uint8_t data[12]; + memcpy(data, &nq->node_id, 8); + data[8] = nq->conn_presence; + data[9] = nq->conn_up; + memcpy(data + 10, &best_rtt, 2); + gui_bridge_post(GUI_EVT_NODEINFO_UPDATE, data, sizeof(data)); +} + #include "gui_bridge_impl.moc" diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index 76dd07aa..260aaf89 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -386,6 +386,8 @@ void UtunNode::runLoop() { etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback); + utun_add_nodeinfo_cbk(m_instance, gui_nodeinfo_cb_impl, nullptr); + /* Bridge chat events to GUI via gui_bridge */ chat_event_set_handler([](int type, const uint8_t* data, int len) { gui_bridge_post(type, data, len);