Browse Source

BGP: фикс владения адресами в реестре (мульти-группа), проброс group_id в сервис (conn_mgr per-group), membership-фильтр create_group; тест media через chat-группу

v2
evgeny 3 weeks ago
parent
commit
42a94aa916
  1. 5
      src/chat/chat_msg.c
  2. 66
      src/media_delivery/media_delivery.c
  3. 6
      src/media_delivery/media_delivery.h
  4. 38
      src/media_delivery/media_download.c
  5. 13
      src/media_delivery/media_index.c
  6. 7
      src/media_delivery/media_index.h
  7. 13
      src/routing_layer/conn_mgr_core.c
  8. 4
      src/routing_layer/conn_mgr_indirect.c
  9. 4
      src/routing_layer/conn_mgr_priv.h
  10. 10
      src/routing_layer/etcp_router.c
  11. 8
      src/routing_layer/etcp_router.h
  12. 17
      src/routing_layer/topo_group.c
  13. 72
      src/routing_layer/topo_node.c
  14. 132
      tests/test_media_delivery_chat.c

5
src/chat/chat_msg.c

@ -13,6 +13,7 @@
#include "../transport_layer/secure_channel.h"
#include "../media_delivery/media_index.h"
#include "../media_delivery/media_download.h"
#include "../media_delivery/media_delivery.h"
#include "../video/video.h"
#include "../../lib/mem.h"
#include "../../lib/platform_compat.h"
@ -193,6 +194,10 @@ static void on_media_registered(void* arg, int err, const struct media_index_res
char attrs[1100]; snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", filename);
chat_core_update_local_attrs(mctx->channel_id, ts, json_sig, attrs);
/* анонсируем блоки суперузлам — иначе первый скачивающий не найдёт держателя */
media_delivery_announce_media(g_cc.inst, strtoull(mctx->channel_id, NULL, 10),
result->media_id, result->block_ids, result->num_blocks);
/* whisper транскрипция для своих голосовых сообщений */
if (mctx->content_type[0] && strncmp(mctx->content_type, "voice", 5) == 0
&& g_cc.inst && g_cc.inst->media_async && g_chat_whisper_trigger) {

66
src/media_delivery/media_delivery.c

@ -9,6 +9,7 @@
#include "../routing_layer/conn_mgr.h"
#include "../transport_layer/etcp_api.h"
#include "../transport_layer/etcp.h"
#include "../transport_layer/secure_channel.h"
#include "../chat/member_sync.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
@ -655,6 +656,8 @@ static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_n
struct media_super_peer* peer = md_super_peer_add(md, from_node, sh->group_id);
if (!peer) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)from_node); return; }
/* группа связи определяется HELLO-пакетом (может отличаться от connect-time, напр. при ручном HELLO) */
peer->group_id = sh->group_id;
peer->peer_last_recv_id = sh->last_recv_id;
peer->connected = 1;
peer->hello_done = 1;
@ -1451,6 +1454,69 @@ int media_delivery_bind(struct UTUN_INSTANCE* inst) {
return media_delivery_init(inst);
}
void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id,
const uint8_t* media_id, const uint8_t* block_ids,
int num_blocks) {
if (!inst || !media_id || !block_ids || num_blocks <= 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce invalid args inst=%p media=%p blocks=%p n=%d",
MD_ID, (void*)inst, (void*)media_id, (void*)block_ids, num_blocks);
return;
}
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) return;
/* собрать суперузлы всех CHAT-групп (node_type=4, не сам) — как md_dl_collect_supernodes */
uint64_t supers[10]; int nsup = 0;
struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0] && inst->topo_sqlite_db) {
char sql[256];
snprintf(sql, sizeof(sql), "SELECT node_id FROM peers_%s WHERE node_type=4 AND node_id!=%llu",
g->channel_id, (unsigned long long)inst->node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) {
while (sqlite3_step(st) == SQLITE_ROW && nsup < 10)
supers[nsup++] = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
}
}
gle = gle->next;
}
if (nsup == 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce — no supernodes for group 0x%016llx",
MD_ID, (unsigned long long)group_id);
return;
}
int64_t ts = (int64_t)time(NULL);
for (int s = 0; s < nsup; s++) {
for (int n = 0; n < num_blocks; n++) {
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK;
hb.group_id = group_id;
memcpy(hb.media_id, media_id, 16);
memcpy(hb.block_id, block_ids + (size_t)n * 16, 16);
hb.chunk = (uint32_t)n;
hb.timestamp = ts;
uint8_t smsg[64]; size_t soff = 0;
memcpy(smsg + soff, hb.block_id, 16); soff += 16;
memcpy(smsg + soff, &hb.chunk, 4); soff += 4;
memcpy(smsg + soff, &hb.timestamp, 8); soff += 8;
uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8;
if (sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce Ed25519 sign failed block=%d", MD_ID, n);
continue;
}
md_send(inst, group_id, supers[s], (const uint8_t*)&hb, sizeof(hb));
}
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce media=%02x%02x... blocks=%d supers=%d",
MD_ID, media_id[0], media_id[1], num_blocks, nsup);
}
void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode) {
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;

6
src/media_delivery/media_delivery.h

@ -130,6 +130,12 @@ int media_delivery_bind(struct UTUN_INSTANCE* inst);
void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode);
void media_delivery_stream_done(struct UTUN_INSTANCE* inst);
/* Автор анонсирует свои блоки суперузлам CHAT-группы (HAVE_BLOCK), чтобы первый
* скачивающий мог найти держателя через QUERY. Вызывается после регистрации медиа. */
void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id,
const uint8_t* media_id, const uint8_t* block_ids,
int num_blocks);
/* relay block context helpers (used by media_download.c) */
struct relay_block_ctx* md_relay_find(struct media_delivery_ctx* md, const uint8_t* block_id);
struct relay_block_ctx* md_relay_add(struct media_delivery_ctx* md,

38
src/media_delivery/media_download.c

@ -639,6 +639,34 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst,
/* ── handle incoming BLOCK_DONE ── */
/* Регистрирует скачанные блоки в media_files (location = собранный файл), чтобы узел
* мог отдавать их как держатель. Вызывается после успешной сборки и проверки хеша. */
static void md_dl_register_servable(struct media_download* dl) {
struct UTUN_INSTANCE* inst = dl->inst;
sqlite3* db = inst ? inst->topo_sqlite_db : NULL;
if (!db) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: register_servable — no DB", MDL_ID); return; }
char chat_id[32];
snprintf(chat_id, sizeof(chat_id), "%llu", (unsigned long long)dl->group_id);
/* location = dest_path относительно media_base (как в media_index_commit) */
const char* rel = dl->dest_path;
size_t base_len = strlen(dl->media_base);
if (strncmp(rel, dl->media_base, base_len) == 0 && rel[base_len] == '/')
rel += base_len + 1;
for (int n = 0; n < dl->num_blocks; n++) {
if (media_index_register_downloaded(db, dl->media_id, dl->blocks[n].block_id,
dl->content_hash, chat_id, rel, inst->node_id,
dl->file_size, dl->blocks[n].expected_size, n,
(int64_t)n * dl->block_size) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: register_servable failed blk=%d", MDL_ID, n);
}
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: registered %d downloaded blocks as servable (%s)",
MDL_ID, dl->num_blocks, rel);
}
void media_download_handle_done(struct UTUN_INSTANCE* inst,
const uint8_t* data, size_t len) {
if (!data || len < MEDIA_BLOCK_DONE_SIZE) return;
@ -757,6 +785,7 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
} }
dl->assembled = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: file assembled: %s", MDL_ID, dl->dest_path);
md_dl_register_servable(dl);
md_dl_finish(dl, 0);
return;
}
@ -906,6 +935,15 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id,
dl = (struct media_download*)qe->data;
/* сохраняем media_base в ui_state — чтобы после скачивания узел мог отдавать блоки
(md_handle_block_req открывает файл по ui_state.media_base + location) */
{
sqlite3_exec(inst->topo_sqlite_db, "CREATE TABLE IF NOT EXISTS ui_state(key TEXT PRIMARY KEY, value TEXT)", NULL, NULL, NULL);
sqlite3_stmt* us = NULL;
sqlite3_prepare_v2(inst->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL);
if (us) { sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); sqlite3_step(us); sqlite3_finalize(us); }
}
md_dl_collect_supernodes(inst, dl);
dl->last_activity_tb = get_time_tb();
dl->watchdog_timer = uasync_set_timeout(inst->ua, (int)md_dl_stall_tb(inst), dl, md_dl_watchdog_cb, "md_dl_wd");

13
src/media_delivery/media_index.c

@ -124,8 +124,7 @@ void media_index_result_free(struct media_index_result* result) {
int media_index_commit(sqlite3* db, const struct media_index_result* result,
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id, const char* dest_path, const char* media_base) {
if (!db || !result || !chat_id || !dest_path || !media_base) return -1;
const char* chat_id, const char* dest_path, const char* media_base) { if (!db || !result || !chat_id || !dest_path || !media_base) return -1;
int nb = result->num_blocks;
int64_t bs = result->block_size;
@ -174,6 +173,16 @@ int media_index_commit(sqlite3* db, const struct media_index_result* result,
return 0;
}
int media_index_register_downloaded(sqlite3* db, const uint8_t* media_id, const uint8_t* block_id,
const uint8_t* content_hash, const char* chat_id,
const char* location, uint64_t node_id,
int64_t file_size, int64_t chunk_size, int chunk, int64_t offset) {
if (!db || !media_id || !block_id || !content_hash || !chat_id || !location) return -1;
static const uint8_t zero_sign[64] = {0};
return mi_insert(db, media_id, block_id, content_hash, chat_id, location,
node_id, zero_sign, (int64_t)time(NULL), file_size, chunk_size, chunk, offset);
}
/* ─── register_async (worker + uasync callback) ─── */
struct mi_reg_ctx {

7
src/media_delivery/media_index.h

@ -33,6 +33,13 @@ int media_index_commit(
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id, const char* dest_path, const char* media_base);
/* Зарегистрировать скачанный блок в media_files (location — путь к собранному файлу),
* чтобы узел мог отдавать его как держатель. node_sign = нули (при отдаче не проверяется). */
int media_index_register_downloaded(sqlite3* db, const uint8_t* media_id, const uint8_t* block_id,
const uint8_t* content_hash, const char* chat_id,
const char* location, uint64_t node_id,
int64_t file_size, int64_t chunk_size, int chunk, int64_t offset);
void media_index_register_async(
struct media_async* ma, struct UASYNC* ua, sqlite3* db,
uint64_t node_id, const uint8_t* ed25519_privkey,

13
src/routing_layer/conn_mgr_core.c

@ -703,8 +703,8 @@ void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
/* Принимающая сторона REVERSE: получили DIRECT_REQ от инициатора — открываем
* NCD-соединение к нему, добавляем линки по адресам из запроса. При INIT
* вызывается cm_reverse_init_cb. */
void cm_handle_direct_req(struct ETCP_CONN* conn, const uint8_t* data, size_t len) {
struct CONN_MGR* mgr = conn->instance->conn_mgr; if (!mgr) return;
void cm_handle_direct_req(struct ETCP_CONN* conn, struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
if (!mgr) return;
struct CM_DIRECT_REQ* req = (struct CM_DIRECT_REQ*)data;
if (len < CM_DIRECT_REQ_H_SIZE + (size_t)req->addr_count * 8) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: DIRECT_REQ too short"); return; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: got DIRECT_REQ from 0x%016llx with %u addrs",
@ -744,12 +744,15 @@ void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry
return;
}
uint8_t* d = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF; uint8_t sub = d[1];
struct CONN_MGR* mgr = conn->instance->conn_mgr;
/* conn_mgr — per-group: резолвим группу из заголовка доставки, а не instance->conn_mgr */
uint64_t group_id; memcpy(&group_id, entry->dgram + ROUTER_SVC_GROUP_OFF, 8);
struct TOPO_GROUP* grp = topo_groups_find(conn->instance->topo_groups, group_id);
struct CONN_MGR* mgr = grp ? grp->conn_mgr : conn->instance->conn_mgr;
if (mgr) { struct CONN_MGR_ENTRY* e = cm_find_entry(mgr, conn->peer_node_id); if (e && e->state == CONN_MGR_STATE_CONNECTED) e->last_traffic_tb = get_time_tb(); }
switch (sub) {
case CM_SUBCMD_DIRECT_REQ: cm_handle_direct_req(conn, d, len); break;
case CM_SUBCMD_DIRECT_REQ: cm_handle_direct_req(conn, mgr, d, len); break;
case CM_SUBCMD_DIRECT_RESP: break;
case CM_SUBCMD_INTERM_EXCHANGE_REQ: if (len >= CM_EXCHANGE_REQ_SIZE) cm_handle_interm_exchange_req(conn, (struct CM_EXCHANGE_REQ*)d); break;
case CM_SUBCMD_INTERM_EXCHANGE_REQ: if (len >= CM_EXCHANGE_REQ_SIZE) cm_handle_interm_exchange_req(conn, mgr, (struct CM_EXCHANGE_REQ*)d); break;
case CM_SUBCMD_INTERM_EXCHANGE_RESP: if (mgr) cm_handle_interm_exchange_resp(mgr, d, len); break;
case CM_SUBCMD_INTERM_SELECTED: if (mgr) cm_handle_interm_selected(mgr, d, len); break;
case CM_SUBCMD_DISCONNECT: if (len >= CM_DISCONNECT_SIZE && mgr) cm_handle_disconnect(mgr, ((struct CM_DISCONNECT*)d)->node_id); break;

4
src/routing_layer/conn_mgr_indirect.c

@ -121,8 +121,8 @@ void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t
/* Цель получила EXCHANGE_REQ: измеряет RTT до кандидатов инициатора (если
* замеры протухли — запускает probe), отправляет EXCHANGE_RESP со СВОИМИ
* кандидатами + замерами RTT до кандидатов инициатора. */
void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CM_EXCHANGE_REQ* req) {
struct CONN_MGR* mgr=conn->instance->conn_mgr; if(!mgr)return;
void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CONN_MGR* mgr, struct CM_EXCHANGE_REQ* req) {
if (!mgr) return;
uint64_t now=get_time_tb();
for(uint8_t i=0;i<req->candidate_count&&i<4;i++){
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,req->candidates[i].node_id);

4
src/routing_layer/conn_mgr_priv.h

@ -142,7 +142,7 @@ void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry);
void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry);
void cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry);
void cm_start_local_scan(struct CONN_MGR_ENTRY* entry);
void cm_handle_direct_req(struct ETCP_CONN* conn, const uint8_t* data, size_t len);
void cm_handle_direct_req(struct ETCP_CONN* conn, struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id);
void cm_reverse_init_cb(struct ETCP_CONN* conn, int event, void* arg);
void cm_reverse_timeout_cb(void* arg);
@ -164,7 +164,7 @@ uint8_t cm_sock_v6_classify(const struct ETCP_SOCKET* s);
void cm_exchange_timeout_cb(void* arg);
void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CM_EXCHANGE_RESP* resp);
void cm_exchange_probe_retry_cb(void* arg);
void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CM_EXCHANGE_REQ* req);
void cm_handle_interm_exchange_req(struct ETCP_CONN* conn, struct CONN_MGR* mgr, struct CM_EXCHANGE_REQ* req);
void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len);

10
src/routing_layer/etcp_router.c

@ -329,6 +329,7 @@ static void router_send_close_to_service(struct ETCP_ROUTER_CONN* rconn) {
memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &rconn->remote_node_id, 8);
memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8);
e->dgram[ROUTER_SVC_FLAGS_OFF] = 0;
memcpy(e->dgram + ROUTER_SVC_GROUP_OFF, &rconn->group_id, 8);
e->len = ROUTER_SVC_HDR_SIZE;
DEBUG_INFO(DEBUG_CATEGORY_ETCPROUTE, "router_close_svc: svc_id=%u remote=%016llx",
rconn->svc_id, (unsigned long long)rconn->remote_node_id);
@ -847,9 +848,9 @@ static void router_idle_ack_timer_cb(void* arg) {
// Reorder / assembly (аналог etcp_output_try_assembly)
// ====================================================================
// Доставка сервисной кодограммы в едином формате [svc_id][src][dst][rx_flags][payload].
// Доставка сервисной кодограммы в едином формате [svc_id][src][dst][rx_flags][group_id][payload].
// Для loopback (dst == self) src и dst оба = self, rx_flags = 0.
static void router_deliver_loopback(struct UTUN_INSTANCE* inst, uint8_t svc_id, struct ll_entry* entry) {
static void router_deliver_loopback(struct UTUN_INSTANCE* inst, uint64_t group_id, uint8_t svc_id, struct ll_entry* entry) {
etcp_recv_fn cb = inst->router_bindings.callbacks[svc_id];
if (!cb) { queue_dgram_free(entry); queue_entry_free(entry); return; }
size_t payload_len = entry->len - 1;
@ -861,6 +862,7 @@ static void router_deliver_loopback(struct UTUN_INSTANCE* inst, uint8_t svc_id,
memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &inst->node_id, 8);
memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8);
e->dgram[ROUTER_SVC_FLAGS_OFF] = 0;
memcpy(e->dgram + ROUTER_SVC_GROUP_OFF, &group_id, 8);
if (payload_len > 0) memcpy(e->dgram + ROUTER_SVC_PAYLOAD_OFF, entry->dgram + 1, payload_len);
e->len = ROUTER_SVC_HDR_SIZE + payload_len;
queue_dgram_free(entry); queue_entry_free(entry);
@ -904,6 +906,7 @@ static void router_deliver(struct ETCP_ROUTER_CONN* rconn, struct ETCP_CONN* con
memcpy(svc_entry->dgram + ROUTER_SVC_SRC_OFF, &rconn->remote_node_id, 8);
memcpy(svc_entry->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8);
svc_entry->dgram[ROUTER_SVC_FLAGS_OFF] = rx_flags;
memcpy(svc_entry->dgram + ROUTER_SVC_GROUP_OFF, &rconn->group_id, 8);
if (payload_len > 0) memcpy(svc_entry->dgram + ROUTER_SVC_PAYLOAD_OFF, payload, payload_len);
if (decoded) u_free(decoded);
@ -1295,7 +1298,7 @@ int etcp_route_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_
if (dst_node_id == inst->node_id) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCPROUTE, "loopback svc_id=%u len=%zu", svc_id, payload_len);
router_deliver_loopback(inst, svc_id, entry);
router_deliver_loopback(inst, group_id, svc_id, entry);
return 0;
}
@ -1324,6 +1327,7 @@ void etcp_router_conn_restart(struct UTUN_INSTANCE* inst, uint64_t group_id, uin
memcpy(e->dgram + ROUTER_SVC_SRC_OFF, &remote_node_id, 8);
memcpy(e->dgram + ROUTER_SVC_DST_OFF, &inst->node_id, 8);
e->dgram[ROUTER_SVC_FLAGS_OFF] = 0;
memcpy(e->dgram + ROUTER_SVC_GROUP_OFF, &group_id, 8);
e->len = ROUTER_SVC_HDR_SIZE;
cb(loopback_conn(inst), e);
} else { queue_entry_free(e); }

8
src/routing_layer/etcp_router.h

@ -40,14 +40,16 @@ struct SVC_ROUTE_HDR {
#define SVC_ROUTE_MAX_BINDINGS 256
// Единый формат доставки сервисной кодограммы (router → сервис):
// [svc_id:1][src_node_id:8][dst_node_id:8][rx_flags:1][payload...]
// [svc_id:1][src_node_id:8][dst_node_id:8][rx_flags:1][group_id:8][payload...]
// src/dst — реальные end-to-end узлы из SVC_ROUTE заголовка (не промежуточные).
// rx_flags — биты ROUTER_FLAG_ENCRYPTED/SIGNED, бывшие на wire (статус для сервиса).
// group_id — группа SVC_ROUTE-пакета (для per-group сервисов вроде conn_mgr).
#define ROUTER_SVC_SRC_OFF 1
#define ROUTER_SVC_DST_OFF 9
#define ROUTER_SVC_FLAGS_OFF 17
#define ROUTER_SVC_PAYLOAD_OFF 18
#define ROUTER_SVC_HDR_SIZE 18 // svc_id(1) + src_node_id(8) + dst_node_id(8) + rx_flags(1)
#define ROUTER_SVC_GROUP_OFF 18
#define ROUTER_SVC_PAYLOAD_OFF 26
#define ROUTER_SVC_HDR_SIZE 26 // svc_id(1) + src_node_id(8) + dst_node_id(8) + rx_flags(1) + group_id(8)
// Биты в flags
#define ROUTER_FLAG_START 0x80

17
src/routing_layer/topo_group.c

@ -511,12 +511,25 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
struct ll_queue* connections = g->instance->connections;
if (connections) {
sqlite3* db = g->instance->topo_sqlite_db;
for (uint32_t slot = 0; slot < connections->hash_size; slot++) {
struct ll_entry* entry = connections->hash_table[slot];
while (entry) {
struct conn_queue_entry* cqe = (struct conn_queue_entry*)entry;
if (cqe->conn && cqe->peer_node_id != 0)
topo_group_new_conn(group, cqe->conn);
if (cqe->conn && cqe->peer_node_id != 0) {
if (group_type == TOPO_GROUP_TYPE_CHAT && db) {
/* CHAT: добавляем conn только если пир — член именно этого канала */
uint64_t* chs = NULL; int chn = 0;
if (topo_node_sqlite_get_member_channels(db, cqe->peer_node_id, &chs, &chn) == 0) {
for (int i = 0; i < chn; i++) {
if (chs[i] == group_id) { topo_group_new_conn(group, cqe->conn); break; }
}
}
u_free(chs);
} else {
topo_group_new_conn(group, cqe->conn);
}
}
entry = entry->hash_next;
}
}

72
src/routing_layer/topo_node.c

@ -132,6 +132,21 @@ uint64_t* topo_node_best_hop_list(struct TOPO_GROUP_NODE* nq, uint8_t* out_count
return (uint64_t*)((uint8_t*)best + sizeof(struct TOPO_NODEPATH));
}
/* Освобождает списки адресов (sock_meta + addrs + reality) узла в глобальном реестре.
* Адреса — глобальная идентичность узла, шарится между группами через node_registry,
* поэтому освобождаются только при последнем unref (topo_node_registry_unref). */
static void topo_node_identity_free(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni) {
if (!groups || !ni) return;
free_v4_sock_list(groups->v4_sock_meta_pool, ni->v4_sock_meta);
free_v4_addr_list(groups->v4_addr_pool, ni->v4_addrs);
free_v6_sock_list(groups->v6_sock_meta_pool, ni->v6_sock_meta);
free_v6_addr_list(groups->v6_addr_pool, ni->v6_addrs);
free_reality_sock_list(ni->reality_socks);
ni->v4_sock_meta = NULL; ni->v4_addrs = NULL;
ni->v6_sock_meta = NULL; ni->v6_addrs = NULL;
ni->reality_socks = NULL;
}
void topo_node_free_raw(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni) {
if (!ni) return;
if (groups) {
@ -151,18 +166,7 @@ void topo_node_destroy(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni) {
void topo_nodeq_free_group_fields(struct TOPO_GROUPS* groups, struct TOPO_GROUP_NODE* nq) {
if (!groups || !nq) return;
struct TOPO_NODE* ni = topo_node_registry_find(groups, nq->node_id);
if (ni) {
free_v4_sock_list(groups->v4_sock_meta_pool, ni->v4_sock_meta);
free_v4_addr_list(groups->v4_addr_pool, ni->v4_addrs);
free_v6_sock_list(groups->v6_sock_meta_pool, ni->v6_sock_meta);
free_v6_addr_list(groups->v6_addr_pool, ni->v6_addrs);
free_reality_sock_list(ni->reality_socks);
ni->v4_sock_meta = NULL; ni->v4_addrs = NULL;
ni->v6_sock_meta = NULL; ni->v6_addrs = NULL;
ni->reality_socks = NULL;
topo_node_registry_unref(groups, nq->node_id);
}
topo_node_registry_unref(groups, nq->node_id);
if (nq->subnets) {
free_v4_sub_list(groups->v4_subnet_pool, nq->subnets->v4_subnets);
free_v6_sub_list(groups->v6_subnet_pool, nq->subnets->v6_subnets);
@ -183,22 +187,42 @@ struct TOPO_NODE* topo_node_registry_find(struct TOPO_GROUPS* groups, uint64_t n
return node;
}
/* Является ли входящий узел полной идентичностью (BGP NODEINFO / локальный узел).
* Partial-узлы из invite/DB-load не несут ed25519 — по ним идентичность не перезаписываем. */
static int topo_node_is_full_identity(const struct TOPO_NODE* ni) {
for (int i = 0; i < SC_PUBKEY_SIZE; i++) if (ni->ed25519_public_key[i]) return 1;
return 0;
}
struct TOPO_NODE* topo_node_registry_store(struct TOPO_GROUPS* groups, struct TOPO_NODE* ni) {
if (!groups || !ni) return NULL;
struct TOPO_NODE* existing = topo_node_registry_find(groups, ni->node_id);
if (existing) {
int n4 = topo_list_count((struct _topo_head*)ni->v4_addrs);
int n6 = topo_list_count((struct _topo_head*)ni->v6_addrs);
int e4 = topo_list_count((struct _topo_head*)existing->v4_addrs);
int e6 = topo_list_count((struct _topo_head*)existing->v6_addrs);
int drop4 = (e4 && n4), drop6 = (e6 && n6);
if (drop4 || drop6)
DEBUG_WARN(DEBUG_CATEGORY_BGP, "registry_store: node=%016llx EXISTS — dropping new addrs v4=%d v6=%d (existing v4=%d v6=%d)",
(unsigned long long)ni->node_id, drop4 ? n4 : 0, drop6 ? n6 : 0, e4, e6);
topo_node_ref(existing);
if (!existing->v4_addrs && ni->v4_addrs) { existing->v4_addrs = ni->v4_addrs; ni->v4_addrs = NULL; }
if (!existing->v6_addrs && ni->v6_addrs) { existing->v6_addrs = ni->v6_addrs; ni->v6_addrs = NULL; }
if (!existing->reality_socks && ni->reality_socks) { existing->reality_socks = ni->reality_socks; ni->reality_socks = NULL; }
/* Идентичность (ver/имя/pubkeys/подпись) — авторитетна у полного NODEINFO,
* иначе реестр держит устаревший ver/имя и форвард nodeinfo ломает подпись. */
if (topo_node_is_full_identity(ni)) {
existing->ver = ni->ver;
existing->client_type = ni->client_type;
existing->client_activity = ni->client_activity;
memcpy(existing->public_key, ni->public_key, SC_PUBKEY_SIZE);
memcpy(existing->ed25519_public_key, ni->ed25519_public_key, SC_PUBKEY_SIZE);
memcpy(existing->x25519_self_sig, ni->x25519_self_sig, 64);
if (ni->node_name) {
if (existing->node_name) u_free(existing->node_name);
existing->node_name = ni->node_name; ni->node_name = NULL;
}
}
/* Адреса — глобальная идентичность. Обновляем каждый список, который входящий узел несёт
* (частичные узлы заполняют только то, что знают, не затирая остальное). */
if (ni->v4_sock_meta) { free_v4_sock_list(groups->v4_sock_meta_pool, existing->v4_sock_meta); existing->v4_sock_meta = ni->v4_sock_meta; ni->v4_sock_meta = NULL; }
if (ni->v4_addrs) { free_v4_addr_list(groups->v4_addr_pool, existing->v4_addrs); existing->v4_addrs = ni->v4_addrs; ni->v4_addrs = NULL; }
if (ni->v6_sock_meta) { free_v6_sock_list(groups->v6_sock_meta_pool, existing->v6_sock_meta); existing->v6_sock_meta = ni->v6_sock_meta; ni->v6_sock_meta = NULL; }
if (ni->v6_addrs) { free_v6_addr_list(groups->v6_addr_pool, existing->v6_addrs); existing->v6_addrs = ni->v6_addrs; ni->v6_addrs = NULL; }
if (ni->reality_socks) { free_reality_sock_list(existing->reality_socks); existing->reality_socks = ni->reality_socks; ni->reality_socks = NULL; }
topo_node_free_raw(groups, ni);
return existing;
}
@ -223,6 +247,8 @@ void topo_node_registry_unref(struct TOPO_GROUPS* groups, uint64_t node_id) {
struct TOPO_NODE* ni = topo_node_registry_find(groups, node_id);
if (!ni) return;
uint32_t prev_ref = ni->group_ref_count;
if (prev_ref == 1)
topo_node_identity_free(groups, ni);
topo_node_unref(ni);
if (prev_ref == 1) {
struct ll_entry* e = groups->node_registry ? queue_find_data_by_index(groups->node_registry, &node_id) : NULL;

132
tests/test_media_delivery_chat.c

@ -136,7 +136,7 @@ static void mon(void* arg) {
if (ce->conn && ce->conn->links) { struct ETCP_LINK* l; for (l = ce->conn->links; l; l = l->next)
if (l->initialized && ce->conn->crypto_ctx.initialized) ok++; } e = e->next; }
}
if (ok >= N_NODES * 2 - 2) g_connected = 1;
if (ok >= N_NODES * (N_NODES - 1)) g_connected = 1; /* полный mesh: N*(N-1) линков */
if (!g_result) uasync_set_timeout(g_ua, 10, NULL, mon, "mdc_mon");
}
@ -162,6 +162,28 @@ static int wait_pubkeys(void) {
return 0;
}
/* ждём, пока BGP CHAT-группы сойдутся: каждый узел имеет маршрут (path) до каждого
* другого мембера в своей CHAT-группе. Без этого media-доставка «через группу» не имеет
* маршрута (topo_group_find_conn_for_node == NULL) и встаёт с «no route». */
static int wait_chat_bgp(void) {
int a = 0;
while (a < 8000) {
int ok = 1;
for (int i = 0; i < N_NODES && ok; i++) {
struct TOPO_GROUP* g = topo_groups_find(g_inst[i]->topo_groups, g_group_id);
if (!g) { ok = 0; break; }
for (int m = 0; m < N_NODES; m++) {
if (m == i) continue;
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, g_nid[m]);
if (!nq || !nq->paths || !nq->paths->head) { ok = 0; break; }
}
}
if (ok) return 1;
uasync_poll(g_ua, POLL_MS); a++;
}
return 0;
}
static void sha256_buf(const uint8_t* data, size_t len, uint8_t out[32]) {
EVP_MD_CTX* ctx = EVP_MD_CTX_new();
EVP_DigestInit_ex(ctx, EVP_sha256(), NULL);
@ -193,6 +215,16 @@ static void ensure_supernode_type(int idx) {
sqlite3_exec(g_inst[idx]->topo_sqlite_db, sql, NULL, NULL, NULL);
}
/* прямой send media-пакета через CHAT-группу (доставка идёт по группе канала) */
static int msend(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, size_t len) {
struct ll_entry* e = queue_entry_new(0); if (!e) return -1;
e->dgram = u_malloc(len + 1); if (!e->dgram) { queue_entry_free(e); return -1; }
e->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY; memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1);
int rc = etcp_route_send(inst, g_group_id, dst, e, 1, 0);
if (rc != 0) { u_free(e->dgram); queue_entry_free(e); }
return rc;
}
/* ── завершение загрузки ── */
static void dl_done_cb(void* arg, int err) {
@ -317,6 +349,14 @@ static int author_index_file(void) {
int rc = media_index_commit(g_inst[I_N1]->topo_sqlite_db, &result, g_nid[I_N1],
g_inst[I_N1]->my_ed25519_privkey, CH_ID, dst_path, media_base);
/* ui_state.media_base — откуда block_req-хендлер открывает файл на стриминг */
{
sqlite3_exec(g_inst[I_N1]->topo_sqlite_db, "CREATE TABLE IF NOT EXISTS ui_state(key TEXT PRIMARY KEY, value TEXT)", NULL, NULL, NULL);
sqlite3_stmt* us = NULL;
sqlite3_prepare_v2(g_inst[I_N1]->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL);
sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); sqlite3_step(us); sqlite3_finalize(us);
}
/* сохраняем глобально */
memcpy(g_media_id, result.media_id, 16);
for (int n = 0; n < MEDIA_NUM_BLOCKS; n++) {
@ -356,7 +396,13 @@ static int author_publish_message(void) {
if (sc_ed25519_sign(g_inst[I_N1]->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) return -1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "[mdc] publish media message ts=%llu len=%zu", (unsigned long long)ts, jl);
return db_sync_insert_signed(g_si[I_N1], json, jl, sig, 64, ts, NULL);
int rc = db_sync_insert_signed(g_si[I_N1], json, jl, sig, 64, ts, NULL);
if (rc == 0) {
/* анонсируем блоки суперузлам (как on_media_registered в chat_msg.c) */
ensure_supernode_type(I_N1);
media_delivery_announce_media(g_inst[I_N1], g_group_id, g_media_id, &g_block_ids[0][0], MEDIA_NUM_BLOCKS);
}
return rc;
}
/* ── фазы ── */
@ -431,6 +477,21 @@ static void setup_chat_group(void) {
media_delivery_set_supernode(g_inst[I_S2], 1);
if (g_inst[I_S1]->md.is_supernode && g_inst[I_S2]->md.is_supernode) OK(); else FAIL();
}
TEST("SUPER_HELLO s1<->s2 (via CHAT)"); {
struct media_pkt_super_hello h; memset(&h, 0, sizeof(h));
h.subcmd = MEDIA_SUBCMD_SUPER_HELLO; h.group_id = g_group_id;
msend(g_inst[I_S1], g_nid[I_S2], (const uint8_t*)&h, sizeof(h));
msend(g_inst[I_S2], g_nid[I_S1], (const uint8_t*)&h, sizeof(h));
int a = 0, ok = 0;
while (a < 1000) {
ok = (g_inst[I_S1]->md.super_peers && g_inst[I_S1]->md.super_peers->head)
&& (g_inst[I_S2]->md.super_peers && g_inst[I_S2]->md.super_peers->head);
if (ok) break;
uasync_poll(g_ua, POLL_MS); a++;
}
if (ok) OK(); else FAIL();
}
}
static void setup_db_sync(void) {
@ -455,10 +516,11 @@ static void setup_db_sync(void) {
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
if (getenv("UTUN_TEST_DEBUG")) {
debug_set_category_level_by_name("media", "trace");
debug_set_category_level_by_name("media", "info");
debug_set_category_level_by_name("general", "info");
debug_set_category_level_by_name("chatsync", "info");
debug_set_category_level_by_name("etcproute", "info");
debug_set_category_level_by_name("chat_sync", "info");
debug_set_category_level_by_name("etcp_route", "info");
debug_set_category_level_by_name("bgp", "trace");
}
utun_instance_set_tun_init_enabled(0);
srand((unsigned)time(NULL));
@ -486,19 +548,17 @@ int main(void) {
for (int i = 0; i < N_NODES; i++) {
char *pr = gv(g_cfg[i], "priv"), *pu = gv(g_cfg[i], "pub");
if (i == 0) {
wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.98.%d.1/24\ntun_ifname=tun%d0\n"
"db_path=%s\ndb_sync_enabled=1\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n",
pr, pu, i, i, g_db_dir[i], g_port[i]);
} else {
int prev = i - 1; char *pv_pu = gv(g_cfg[prev], "pub");
char clients[4096] = "";
for (int j = 0; j < i; j++) {
char *pj_pu = gv(g_cfg[j], "pub");
char link[256]; snprintf(link, sizeof(link),
"[client: to_n%d]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", prev, pv_pu, g_port[prev]);
wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.98.%d.1/24\ntun_ifname=tun%d0\n"
"db_path=%s\ndb_sync_enabled=1\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n%s[allowed_keys]\nallow_all=1\n",
pr, pu, i, i, g_db_dir[i], g_port[i], link);
u_free(pv_pu);
"[client: to_n%d]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", j, pj_pu, g_port[j]);
strncat(clients, link, sizeof(clients) - strlen(clients) - 1);
u_free(pj_pu);
}
wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.98.%d.1/24\ntun_ifname=tun%d0\n"
"db_path=%s\ndb_sync_enabled=1\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n%s[allowed_keys]\nallow_all=1\n",
pr, pu, i, i, g_db_dir[i], g_port[i], clients);
u_free(pr); u_free(pu);
}
@ -506,6 +566,7 @@ int main(void) {
g_inst[i] = utun_instance_create(g_ua, g_cfg[i]);
if (!g_inst[i]) { printf("FAIL: create n%d\n", i); goto done; }
utun_instance_init(g_inst[i]);
member_sync_init(g_inst[i]);
}
mon(NULL);
@ -516,6 +577,10 @@ int main(void) {
g_group_id = strtoull(CH_ID, NULL, 10);
/* media_files-таблица нужна на КАЖДОМ узле (в chat_core_init её создаёт media_index_init) —
иначе узел, принявший BLOCK_REQ, падает с "no such table: media_files". */
for (int i = 0; i < N_NODES; i++) media_index_init(g_inst[i]->topo_sqlite_db);
setup_chat_group();
setup_db_sync();
@ -523,6 +588,10 @@ int main(void) {
if (wait_pubkeys()) OK(); else FAIL();
}
TEST("CHAT group BGP converged (routing ready)"); {
if (wait_chat_bgp()) OK(); else FAIL();
}
TEST("author: index file + media_index_commit"); {
if (author_index_file() == 0 && db_count(g_inst[I_N1]->topo_sqlite_db, "media_files", NULL, 0) == MEDIA_NUM_BLOCKS) OK();
else FAIL();
@ -534,12 +603,13 @@ int main(void) {
/* Фаза 1: сообщение распространилось на все узлы */
TEST("message propagated to all 6 nodes"); {
int a = 0, all = 0;
while (a < 4000) {
int all = 0;
uint64_t start = get_time_tb();
while ((get_time_tb() - start) < (uint64_t)60000) {
all = 1;
for (int i = 0; i < N_NODES; i++) if (db_sync_count(g_si[i]) < 1) all = 0;
if (all) break;
uasync_poll(g_ua, POLL_MS); a++;
uasync_poll(g_ua, POLL_MS);
}
if (all) OK(); else FAIL("counts: %d %d %d %d %d %d",
db_sync_count(g_si[0]), db_sync_count(g_si[1]), db_sync_count(g_si[2]),
@ -550,8 +620,9 @@ int main(void) {
TEST("storage st1/st2 downloaded + assembled media"); {
char d1[1024], d2[1024];
node_dest(I_ST1, d1, sizeof(d1)); node_dest(I_ST2, d2, sizeof(d2));
int a = 0;
while (a < 8000 && (!g_dl_done[I_ST1] || !g_dl_done[I_ST2])) { uasync_poll(g_ua, POLL_MS); a++; }
uint64_t start = get_time_tb();
while ((get_time_tb() - start) < (uint64_t)60000 && (!g_dl_done[I_ST1] || !g_dl_done[I_ST2]))
uasync_poll(g_ua, POLL_MS);
if (g_dl_done[I_ST1] && g_dl_done[I_ST2] && g_dl_err[I_ST1] == 0 && g_dl_err[I_ST2] == 0
&& file_exists(d1) && file_exists(d2)) OK();
else FAIL("st1 done=%d err=%d exists=%d; st2 done=%d err=%d exists=%d",
@ -569,12 +640,13 @@ int main(void) {
/* Фаза 3: суперузлы узнали о блоках + репликация между собой */
TEST("supernodes have block_availability (from storage HAVE_BLOCK)"); {
int a = 0, ok = 0;
while (a < 2000) {
int ok = 0;
uint64_t start = get_time_tb();
while ((get_time_tb() - start) < (uint64_t)60000) {
int n1 = db_count(g_inst[I_S1]->topo_sqlite_db, "block_availability", NULL, 0);
int n2 = db_count(g_inst[I_S2]->topo_sqlite_db, "block_availability", NULL, 0);
if (n1 >= 2 && n2 >= 2) { ok = 1; break; }
uasync_poll(g_ua, POLL_MS); a++;
uasync_poll(g_ua, POLL_MS);
}
if (ok) OK(); else FAIL("s1=%d s2=%d",
db_count(g_inst[I_S1]->topo_sqlite_db, "block_availability", NULL, 0),
@ -582,15 +654,16 @@ int main(void) {
}
TEST("supernode replication s1↔s2 (super_sync)"); {
int a = 0, ok = 0;
while (a < 2000) {
int ok = 0;
uint64_t start = get_time_tb();
while ((get_time_tb() - start) < (uint64_t)60000) {
sqlite3_stmt* st = NULL; uint64_t lr = 0;
sqlite3_prepare_v2(g_inst[I_S2]->topo_sqlite_db, "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?", -1, &st, NULL);
sqlite3_bind_int64(st, 1, (sqlite3_int64)g_nid[I_S1]);
if (sqlite3_step(st) == SQLITE_ROW) lr = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
if (lr > 0) { ok = 1; break; }
uasync_poll(g_ua, POLL_MS); a++;
uasync_poll(g_ua, POLL_MS);
}
if (ok) OK(); else FAIL();
}
@ -612,8 +685,9 @@ int main(void) {
media_download_start(g_inst[I_N2], g_group_id, &r, dest, mbase, g_nid[I_N1], dl_done_cb, (void*)(intptr_t)I_N2, NULL, NULL);
media_index_result_free(&r);
int a = 0;
while (a < 8000 && !g_dl_done[I_N2]) { uasync_poll(g_ua, POLL_MS); a++; }
uint64_t start = get_time_tb();
while ((get_time_tb() - start) < (uint64_t)60000 && !g_dl_done[I_N2])
uasync_poll(g_ua, POLL_MS);
if (g_dl_done[I_N2] && g_dl_err[I_N2] == 0 && file_exists(dest)) OK();
else FAIL("done=%d err=%d exists=%d", g_dl_done[I_N2], g_dl_err[I_N2], file_exists(dest));
}

Loading…
Cancel
Save