From 8a9184675a3a0bcafc1d3e1451ae6700222163fc Mon Sep 17 00:00:00 2001 From: evgeny Date: Sat, 19 Sep 2026 18:16:42 +0300 Subject: [PATCH] =?UTF-8?q?radio:=20=D0=B4=D0=B5=D0=B4=D1=83=D0=BF=D0=BB?= =?UTF-8?q?=D0=B8=D0=BA=D0=B0=D1=86=D0=B8=D1=8F=20=D0=B4=D1=83=D0=B1=D0=BB?= =?UTF-8?q?=D0=B5=D0=B9=20=D0=B2=20mesh=20+=20=D1=84=D0=B8=D0=BA=D1=81=20O?= =?UTF-8?q?OB=20=D0=B2=20hop=5Flist?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - radio_recv: кэш (src,stream,seq) с TTL 2с гасит кадры, пришедшие по второму пути (diamond/mesh) - фикс записи за границу new_hops[hop_count] при hop_count==RADIO_MAX_HOPS - счётчик c_dup_dropped + radio_dup_dropped() для диагностики - тест test_radio_mesh (4 узла, diamond): C получает ровно N кадров, dup>0 --- src/radio/radio.c | 63 +++++- src/radio/radio.h | 3 + tests/Makefile.am | 5 + tests/test_radio_mesh.c | 474 ++++++++++++++++++++++++++++++++++++++++ 4 files changed, 539 insertions(+), 6 deletions(-) create mode 100644 tests/test_radio_mesh.c diff --git a/src/radio/radio.c b/src/radio/radio.c index 6d34347c..15a18c95 100644 --- a/src/radio/radio.c +++ b/src/radio/radio.c @@ -29,6 +29,17 @@ #define RADIO_MAX_TALKERS 16 /* максимум одновременно говорящих источников */ +#define RADIO_DEDUP_MAX 128 /* ёмкость кэша дедупликации (≈2.5с @ 50 кадров/с) */ +#define RADIO_DEDUP_TTL_TB 20000 /* 2с — окно дедупликации (совпадает с таймаутом источника) */ + +/* Запись кэша дедупликации: ключ — (src, stream, seq), tb — момент прихода. */ +struct radio_seen { + uint64_t src; + uint16_t stream; + uint16_t seq; + uint64_t tb; +}; + struct radio_ctx { struct TOPO_GROUP* group; uint16_t stream_id; /* id текущей передачи (talk-burst) */ @@ -40,7 +51,10 @@ struct radio_ctx { uint32_t c_loop_dropped; /* дропнуто из-за цикла (self в hop_list) */ uint32_t c_nosub_dropped; /* не отправлено: за соседом нет подписчиков */ uint32_t c_bad_dropped; /* дропнуто некорректных пакетов */ + uint32_t c_dup_dropped; /* дропнуто дублей (src,stream,seq) пришедших по другому пути */ struct radio_talker { uint64_t src; uint16_t stream; } talkers[RADIO_MAX_TALKERS]; /* активные говорящие */ + struct radio_seen seen[RADIO_DEDUP_MAX]; /* кэш дедупликации */ + int seen_count; }; struct radio_instance { @@ -73,10 +87,10 @@ int radio_init(struct TOPO_GROUP* group) { void radio_destroy(struct TOPO_GROUP* group) { if (!group || !group->radio) return; struct radio_ctx* ctx = group->radio; - DEBUG_INFO(DEBUG_CATEGORY_RADIO, "%s: destroy grp=%016llx (tx=%u rx=%u fwd=%u loop=%u nosub=%u bad=%u)", + DEBUG_INFO(DEBUG_CATEGORY_RADIO, "%s: destroy grp=%016llx (tx=%u rx=%u fwd=%u loop=%u nosub=%u bad=%u dup=%u)", RADIO_ID, (unsigned long long)group->group_id, ctx->c_tx_frames, ctx->c_rx_frames, ctx->c_fwd_frames, - ctx->c_loop_dropped, ctx->c_nosub_dropped, ctx->c_bad_dropped); + ctx->c_loop_dropped, ctx->c_nosub_dropped, ctx->c_bad_dropped, ctx->c_dup_dropped); u_free(ctx); group->radio = NULL; } @@ -207,6 +221,29 @@ static void radio_talker_off(struct radio_ctx* ctx, uint64_t src) { } } +/* Дедупликация: вернуть 1, если (src,stream,seq) уже видели (дубль по другому пути), + * иначе зафиксировать в кэше и вернуть 0. Лениво чистит просроченные записи. */ +static int radio_seen_dup(struct radio_ctx* ctx, uint64_t src, uint16_t stream, uint16_t seq, uint64_t now) { + int free_slot = -1, oldest = -1; + uint64_t oldest_tb = 0; + for (int i = 0; i < ctx->seen_count; i++) { + struct radio_seen* s = &ctx->seen[i]; + if (s->tb && now - s->tb > RADIO_DEDUP_TTL_TB) { s->tb = 0; s->src = 0; } + if (!s->tb) { if (free_slot < 0) free_slot = i; continue; } + if (s->src == src && s->stream == stream && s->seq == seq) return 1; + if (oldest < 0 || s->tb < oldest_tb) { oldest = i; oldest_tb = s->tb; } + } + if (free_slot < 0) { + if (ctx->seen_count < RADIO_DEDUP_MAX) free_slot = ctx->seen_count++; + else free_slot = oldest; + } + ctx->seen[free_slot].src = src; + ctx->seen[free_slot].stream = stream; + ctx->seen[free_slot].seq = seq; + ctx->seen[free_slot].tb = now; + return 0; +} + /* ── приём ── */ static void radio_recv(struct TOPO_GROUP* group, struct ETCP_CONN* from_conn, @@ -248,6 +285,13 @@ static void radio_recv(struct TOPO_GROUP* group, struct ETCP_CONN* from_conn, } } + /* дедупликация: (src,stream,seq) уже приходил по другому пути (mesh) */ + if (radio_seen_dup(ctx, src, stream, seq, get_time_tb())) { + ctx->c_dup_dropped++; + DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: dup drop src=%016llx stream=%u seq=%u", RADIO_ID, (unsigned long long)src, stream, seq); + return; + } + /* доставка локально, если мы подписаны */ if (group->radio_active && ctx->group) { struct radio_instance* ri = radio_inst(group->instance); @@ -262,15 +306,14 @@ static void radio_recv(struct TOPO_GROUP* group, struct ETCP_CONN* from_conn, } /* форвардинг: новый hop_list = старый + self */ - uint64_t new_hops[RADIO_MAX_HOPS]; - memcpy(new_hops, hops, (size_t)hop_count * 8); - new_hops[hop_count] = self; uint8_t new_count = hop_count + 1; - if (new_count >= RADIO_MAX_HOPS) { DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: max hops reached, no fwd src=%016llx", RADIO_ID, (unsigned long long)src); return; } + uint64_t new_hops[RADIO_MAX_HOPS]; + memcpy(new_hops, hops, (size_t)hop_count * 8); + new_hops[hop_count] = self; uint8_t out[RADIO_HDR_SIZE + RADIO_MAX_HOPS * 8 + RADIO_MAX_OPUS]; size_t out_len = radio_build(out, group_id, src, stream, seq, flags, new_hops, new_count, opus, opus_len); @@ -348,6 +391,14 @@ int radio_subscribers(struct UTUN_INSTANCE* inst, uint64_t group_id) { return group ? topo_group_radio_subscribers(group) : 0; } +/* Число дропнутых дублей (src,stream,seq) — для диагностики/тестов. */ +uint32_t radio_dup_dropped(struct UTUN_INSTANCE* inst, uint64_t group_id) { + if (!inst || !inst->topo_groups) return 0; + struct TOPO_GROUP* group = topo_groups_find(inst->topo_groups, group_id); + struct radio_ctx* ctx = radio_of(group); + return ctx ? ctx->c_dup_dropped : 0; +} + /* ── передача (PTT) ── */ int radio_talk_begin(struct UTUN_INSTANCE* inst, uint64_t group_id) { diff --git a/src/radio/radio.h b/src/radio/radio.h index 5734a077..27221600 100644 --- a/src/radio/radio.h +++ b/src/radio/radio.h @@ -44,6 +44,9 @@ int radio_is_active(struct UTUN_INSTANCE* inst, uint64_t group_id); /* Число подписчиков рации канала (без себя) — для индикатора в GUI. */ int radio_subscribers(struct UTUN_INSTANCE* inst, uint64_t group_id); +/* Число дропнутых дублей (src,stream,seq) — для диагностики/тестов. */ +uint32_t radio_dup_dropped(struct UTUN_INSTANCE* inst, uint64_t group_id); + /* ── trampoline-структуры (GUI/аудио-поток → uasync) ── */ struct radio_talk_arg { diff --git a/tests/Makefile.am b/tests/Makefile.am index a2dae42e..768aa187 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -69,6 +69,7 @@ check_PROGRAMS = \ test_call_headless \ test_radio \ test_radio_headless \ + test_radio_mesh \ test_stcp_traffic \ test_bbr_integration \ test_audio_compressor \ @@ -417,6 +418,10 @@ test_radio_headless_SOURCES = test_radio_headless.c test_radio_headless_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/chat -I$(top_srcdir)/src/dm -I$(top_srcdir)/src/call -I$(top_srcdir)/src/radio -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/lib test_radio_headless_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_radio_mesh_SOURCES = test_radio_mesh.c +test_radio_mesh_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/chat -I$(top_srcdir)/src/dm -I$(top_srcdir)/src/call -I$(top_srcdir)/src/radio -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/lib +test_radio_mesh_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_stcp_traffic_SOURCES = test_stcp_traffic.c test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_stcp_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_radio_mesh.c b/tests/test_radio_mesh.c new file mode 100644 index 00000000..1a124a8d --- /dev/null +++ b/tests/test_radio_mesh.c @@ -0,0 +1,474 @@ +// test_radio_mesh.c — интеграционный тест рации: дедупликация дублей в mesh (diamond). +// +// Однопоточный: 4 узла (A, B, D, C) в одном процессе, один UASYNC, без fork. +// Топология — diamond: A↔B, A↔D, B↔C, D↔C (A сервер, B и D — клиенты к A + +// серверы, C — клиент к B и к D). Все четыре в одной CHAT-группе. +// Master state-machine на таймере. +// +// Сценарий: +// 1. A,B,D,C radio_on → флаг TOPO_FLAG_RADIO распространяется по BGP. +// 2. A: radio_talk_begin + N кадров + radio_talk_end. +// Кадр идёт к C двумя путями (A→B→C и A→D→C) — дубль гасится кэшем +// (src,stream,seq) в radio_recv. +// 3. Проверка: B, D, C получают ровно N кадров (без дублей), порядок seq +// и FIN корректны; у C счётчик дропнутых дублей > 0 (путь дедупликации +// реально отработал). + +#include "chat_core.h" +#include "chat_sync.h" +#include "chat_event.h" +#include "member_sync.h" +#include "../radio/radio.h" +#include "../routing_layer/topo_node_sqlite.h" +#include "../routing_layer/topo_group.h" +#include "../routing_layer/topo_node.h" +#include "../utun_instance.h" +#include "../transport_layer/etcp.h" +#include "../transport_layer/secure_channel.h" +#include "../ntp_time.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/platform_compat.h" + +#include +#define OPENSSL_API_COMPAT 0x10100000L +#include +#include +#include +#include +#include +#include + +#define CH_ID "4242424242424242" + +#define IDX_A 0 +#define IDX_B 1 +#define IDX_D 2 +#define IDX_C 3 + +#define TICK_TB 100 /* 10 ms на такт */ +#define MAX_TICKS 4000 /* глобальный таймаут ≈ 40 c */ +#define N_FRAMES 5 + +enum mesh_phase { + P_WAIT_CONN, + P_SETUP, + P_WAIT_BGP, + P_RADIO_ON, + P_WAIT_FLAGS, + P_TALK1, + P_WAIT_RX1, + P_DONE, +}; + +struct mesh_ctx { + struct UASYNC* ua; + struct UTUN_INSTANCE* inst[4]; + enum mesh_phase phase; + int result; /* 0=running, 1=ok, 2=fail */ + int ticks; + void* timer; +}; + +/* приёмник кадров (per-узел): проверка порядка seq и FIN */ +struct rx_evt { + int active; + uint16_t stream; + uint16_t expect_seq; + int frames; /* число data-кадров (fin=0) текущего потока */ + int fin_seen; + int seq_ok; +}; + +struct mesh_shared { + uint8_t ch_x_pub[32], ch_ed_pub[32], ch_ed_priv[32], ch_sig[64]; + uint64_t nid[4]; + uint8_t x_pub[4][32], x_priv[4][32], ed_pub[4][32]; + int port[4]; +}; +static struct mesh_shared g_sh; +static struct rx_evt g_rx[4]; + +static int G_PASSED = 0, G_FAILED = 0, G_TOTAL = 0; +#define TEST(n) do { G_TOTAL++; printf(" %-58s", n); fflush(stdout); } while(0) +#define OK() do { G_PASSED++; printf("OK\n"); } while(0) +#define FAIL(f,...) do { G_FAILED++; printf("FAIL: " f "\n", ##__VA_ARGS__); } while(0) + +/* ── крипто-хелперы ── */ + +static void gen_x25519(uint8_t pub[32], uint8_t priv[32]) { + EVP_PKEY_CTX* ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_X25519, NULL); + EVP_PKEY* pkey = NULL; + EVP_PKEY_keygen_init(ctx); EVP_PKEY_keygen(ctx, &pkey); EVP_PKEY_CTX_free(ctx); + size_t l = 32; EVP_PKEY_get_raw_public_key(pkey, pub, &l); + l = 32; EVP_PKEY_get_raw_private_key(pkey, priv, &l); + EVP_PKEY_free(pkey); +} + +static void gen_ed25519(uint8_t pub[32], uint8_t priv[32]) { + EVP_PKEY_CTX* ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_ED25519, NULL); + EVP_PKEY* pkey = NULL; + EVP_PKEY_keygen_init(ctx); EVP_PKEY_keygen(ctx, &pkey); EVP_PKEY_CTX_free(ctx); + size_t l = 32; EVP_PKEY_get_raw_public_key(pkey, pub, &l); + l = 32; EVP_PKEY_get_raw_private_key(pkey, priv, &l); + EVP_PKEY_free(pkey); +} + +static void to_hex(const uint8_t* bin, size_t n, char* out) { + for (size_t i = 0; i < n; i++) sprintf(out + i * 2, "%02x", bin[i]); + out[n * 2] = '\0'; +} + +static int wf(const char* p, const char* f, ...) { + va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; + va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; +} + +/* ── данные сценария + конфиги (diamond A↔B, A↔D, B↔C, D↔C) ── */ + +static void fill_shared(int base_port) { + memset(&g_sh, 0, sizeof(g_sh)); + for (int i = 0; i < 4; i++) { + gen_x25519(g_sh.x_pub[i], g_sh.x_priv[i]); + sc_derive_ed25519_pubkey(g_sh.x_priv[i], g_sh.ed_pub[i]); + g_sh.nid[i] = sc_derive_node_id_from_pubkey(g_sh.x_pub[i]); + g_sh.port[i] = base_port + i * 1000; + } + { uint8_t ch_x_priv[32]; gen_x25519(g_sh.ch_x_pub, ch_x_priv); } + gen_ed25519(g_sh.ch_ed_pub, g_sh.ch_ed_priv); + { + uint8_t msg[256]; size_t off = 0; + off += (size_t)snprintf((char*)msg + off, sizeof(msg) - off, "%s", CH_ID) + 1; + off += (size_t)snprintf((char*)msg + off, sizeof(msg) - off, "%s", "radio-mesh") + 1; + memcpy(msg + off, &g_sh.nid[IDX_A], 8); off += 8; + memcpy(msg + off, g_sh.ch_x_pub, 32); off += 32; + memcpy(msg + off, g_sh.ch_ed_pub, 32); off += 32; + sc_ed25519_sign(g_sh.ch_ed_priv, msg, off, g_sh.ch_sig); + } +} + +static void write_configs(const char* dir) { + for (int i = 0; i < 4; i++) { + char priv[65], pub[65], path[512], dbd[512]; + to_hex(g_sh.x_priv[i], 32, priv); + to_hex(g_sh.x_pub[i], 32, pub); + snprintf(path, sizeof(path), "%s/%s.conf", dir, + i == IDX_A ? "a" : i == IDX_B ? "b" : i == IDX_D ? "d" : "c"); + snprintf(dbd, sizeof(dbd), "%s/db%s", dir, + i == IDX_A ? "a" : i == IDX_B ? "b" : i == IDX_D ? "d" : "c"); + utun_mkdir(dbd, 0755); + + char client[1024] = ""; + if (i == IDX_B) { /* B → A */ + char apub[65]; to_hex(g_sh.x_pub[IDX_A], 32, apub); + snprintf(client, sizeof(client), "[client: to_a]\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", apub, g_sh.port[IDX_A]); + } else if (i == IDX_D) { /* D → A */ + char apub[65]; to_hex(g_sh.x_pub[IDX_A], 32, apub); + snprintf(client, sizeof(client), "[client: to_a]\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", apub, g_sh.port[IDX_A]); + } else if (i == IDX_C) { /* C → B и C → D */ + char bpub[65], dpub[65]; + to_hex(g_sh.x_pub[IDX_B], 32, bpub); + to_hex(g_sh.x_pub[IDX_D], 32, dpub); + snprintf(client, sizeof(client), + "[client: to_b]\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n" + "[client: to_d]\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", + bpub, g_sh.port[IDX_B], dpub, g_sh.port[IDX_D]); + } + + wf(path, + "[global]\ntun_ip=10.98.%d.1/24\ntun_ifname=tun%d0\ndb_path=%s/db%s\n" + "my_private_key=%s\nmy_public_key=%s\n" + "[server: s1]\naddr=127.0.0.1:%d\ntype=public\n" + "%s" + "[chatserver]\nstorage_autoload=0\n[allowed_keys]\nallow_all=1\n", + i, i, dir, i == IDX_A ? "a" : i == IDX_B ? "b" : i == IDX_D ? "d" : "c", + priv, pub, g_sh.port[i], client); + } +} + +/* ── проверки ── */ + +static int conn_up(struct UTUN_INSTANCE* inst, uint64_t nid) { + struct ETCP_CONN* c = instance_find_conn(inst, nid); + return c && c->initialized && c->links_up; +} + +static int bgp_has(struct UTUN_INSTANCE* inst, uint64_t nid) { + uint64_t gid = strtoull(CH_ID, NULL, 10); + struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid); + if (!g) return 0; + struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, nid); + return nq && nq->paths && nq->paths->head; +} + +static int node_radio(struct UTUN_INSTANCE* inst, uint64_t nid) { + uint64_t gid = strtoull(CH_ID, NULL, 10); + struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, gid); + if (!g) return 0; + struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(g, nid); + return nq ? nq->radio : 0; +} + +/* все 8 направлений diamond подняты */ +static int all_conn_up(struct mesh_ctx* t) { + return conn_up(t->inst[IDX_A], g_sh.nid[IDX_B]) && conn_up(t->inst[IDX_B], g_sh.nid[IDX_A]) + && conn_up(t->inst[IDX_A], g_sh.nid[IDX_D]) && conn_up(t->inst[IDX_D], g_sh.nid[IDX_A]) + && conn_up(t->inst[IDX_B], g_sh.nid[IDX_C]) && conn_up(t->inst[IDX_C], g_sh.nid[IDX_B]) + && conn_up(t->inst[IDX_D], g_sh.nid[IDX_C]) && conn_up(t->inst[IDX_C], g_sh.nid[IDX_D]); +} + +/* каждый узел видит 3 остальных в BGP */ +static int all_bgp(struct mesh_ctx* t) { + for (int i = 0; i < 4; i++) + for (int j = 0; j < 4; j++) + if (i != j && !bgp_has(t->inst[i], g_sh.nid[j])) return 0; + return 1; +} + +/* каждый узел видит radio=1 у 3 остальных */ +static int all_radio(struct mesh_ctx* t) { + for (int i = 0; i < 4; i++) + for (int j = 0; j < 4; j++) + if (i != j && !node_radio(t->inst[i], g_sh.nid[j])) return 0; + return 1; +} + +/* ── frame cb (per-узел) ── */ + +static void rx_frame_cb(struct UTUN_INSTANCE* inst, uint64_t group_id, + uint64_t src, uint16_t stream, uint16_t seq, + uint8_t fin, const uint8_t* opus, int len, void* arg) { + (void)inst; (void)group_id; (void)src; (void)opus; (void)len; + struct rx_evt* rx = (struct rx_evt*)arg; + if (fin) { + if (rx->active && stream == rx->stream) rx->fin_seen = 1; + return; + } + if (!rx->active || stream != rx->stream) { + rx->active = 1; + rx->stream = stream; + rx->expect_seq = 0; + rx->frames = 0; + rx->fin_seen = 0; + rx->seq_ok = 1; + } + if (seq != rx->expect_seq) rx->seq_ok = 0; + rx->expect_seq = (uint16_t)(seq + 1); + rx->frames++; +} + +/* ── настройка канала и мемберов ── */ + +static void channel_put_shared(struct UTUN_INSTANCE* inst, int has_priv) { + topo_node_sqlite_channel_put(inst->topo_sqlite_db, CH_ID, "radio-mesh", g_sh.nid[IDX_A], + g_sh.ch_x_pub, NULL, g_sh.ch_ed_pub, has_priv ? g_sh.ch_ed_priv : NULL, g_sh.ch_sig); + chat_core_ensure_channel_ready(inst, CH_ID); +} + +static void member_put_self(struct UTUN_INSTANCE* inst, const char* name) { + uint64_t join_ts = (uint64_t)ntp_time_get_seconds(inst); + uint8_t jm[128]; + int jl = member_sync_build_join_msg(g_sh.ch_x_pub, g_sh.ch_ed_pub, inst->node_id, + inst->my_keys.public_key, join_ts, jm, (int)sizeof(jm)); + uint8_t join_sig[64]; sc_ed25519_sign(inst->my_ed25519_privkey, jm, (size_t)jl, join_sig); + char userinfo[256]; snprintf(userinfo, sizeof(userinfo), "{\"name\":\"%s\"}", name); + member_sync_put(inst, CH_ID, inst->node_id, inst->my_keys.public_key, inst->my_ed25519_pubkey, + join_sig, join_ts, NULL, 0, userinfo, NULL, NULL, 0, 0, NULL); +} + +static void member_put_placeholder(struct UTUN_INSTANCE* inst, int idx, const char* name) { + char userinfo[256]; snprintf(userinfo, sizeof(userinfo), "{\"name\":\"%s\"}", name); + member_sync_put(inst, CH_ID, g_sh.nid[idx], g_sh.x_pub[idx], g_sh.ed_pub[idx], + NULL, 0, NULL, 0, userinfo, NULL, NULL, 0, 0, NULL); +} + +static void do_setup(struct UTUN_INSTANCE* inst, int role) { + switch (role) { + case IDX_A: + channel_put_shared(inst, 1); + member_put_self(inst, "A"); + member_put_placeholder(inst, IDX_B, "B"); + member_put_placeholder(inst, IDX_D, "D"); + member_put_placeholder(inst, IDX_C, "C"); + member_sync_start(inst, g_sh.nid[IDX_B], CH_ID, NULL, NULL); + member_sync_start(inst, g_sh.nid[IDX_D], CH_ID, NULL, NULL); + break; + case IDX_B: + channel_put_shared(inst, 0); + member_put_self(inst, "B"); + member_put_placeholder(inst, IDX_A, "A"); + member_put_placeholder(inst, IDX_D, "D"); + member_put_placeholder(inst, IDX_C, "C"); + member_sync_start(inst, g_sh.nid[IDX_A], CH_ID, NULL, NULL); + member_sync_start(inst, g_sh.nid[IDX_C], CH_ID, NULL, NULL); + break; + case IDX_D: + channel_put_shared(inst, 0); + member_put_self(inst, "D"); + member_put_placeholder(inst, IDX_A, "A"); + member_put_placeholder(inst, IDX_B, "B"); + member_put_placeholder(inst, IDX_C, "C"); + member_sync_start(inst, g_sh.nid[IDX_A], CH_ID, NULL, NULL); + member_sync_start(inst, g_sh.nid[IDX_C], CH_ID, NULL, NULL); + break; + case IDX_C: + channel_put_shared(inst, 0); + member_put_self(inst, "C"); + member_put_placeholder(inst, IDX_A, "A"); + member_put_placeholder(inst, IDX_B, "B"); + member_put_placeholder(inst, IDX_D, "D"); + member_sync_start(inst, g_sh.nid[IDX_B], CH_ID, NULL, NULL); + member_sync_start(inst, g_sh.nid[IDX_D], CH_ID, NULL, NULL); + break; + } +} + +/* ── talk: N кадров + FIN от A ── */ + +static void do_talk(struct UTUN_INSTANCE* A) { + uint64_t gid = strtoull(CH_ID, NULL, 10); + static const uint8_t opus[16] = { 0xf8, 0xff, 0xfe, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 0x08, 0x09, 0x0a, 0x0b, 0x0c, 0x0d }; + radio_talk_begin(A, gid); + for (int i = 0; i < N_FRAMES; i++) radio_talk_send(A, gid, opus, sizeof(opus)); + radio_talk_end(A, gid); +} + +/* ── master state-machine ── */ + +static void mesh_tick(void* arg) { + struct mesh_ctx* t = arg; + t->timer = NULL; + if (t->result) return; + + struct UTUN_INSTANCE* A = t->inst[IDX_A]; + struct UTUN_INSTANCE* B = t->inst[IDX_B]; + struct UTUN_INSTANCE* D = t->inst[IDX_D]; + struct UTUN_INSTANCE* C = t->inst[IDX_C]; + uint64_t gid = strtoull(CH_ID, NULL, 10); + + switch (t->phase) { + case P_WAIT_CONN: + if (all_conn_up(t)) t->phase = P_SETUP; + break; + case P_SETUP: + do_setup(A, IDX_A); + do_setup(B, IDX_B); + do_setup(D, IDX_D); + do_setup(C, IDX_C); + t->phase = P_WAIT_BGP; + break; + case P_WAIT_BGP: + if (all_bgp(t)) t->phase = P_RADIO_ON; + break; + case P_RADIO_ON: + radio_set_active(A, gid, 1); + radio_set_active(B, gid, 1); + radio_set_active(D, gid, 1); + radio_set_active(C, gid, 1); + t->phase = P_WAIT_FLAGS; + break; + case P_WAIT_FLAGS: + if (all_radio(t)) t->phase = P_TALK1; + break; + case P_TALK1: + memset(g_rx, 0, sizeof(g_rx)); + do_talk(A); + t->phase = P_WAIT_RX1; + break; + case P_WAIT_RX1: + /* B и D получают по N кадров; C — ровно N (дубль по второму пути гасится кэшем) */ + if (g_rx[IDX_B].frames == N_FRAMES && g_rx[IDX_B].fin_seen && g_rx[IDX_B].seq_ok + && g_rx[IDX_D].frames == N_FRAMES && g_rx[IDX_D].fin_seen && g_rx[IDX_D].seq_ok + && g_rx[IDX_C].frames == N_FRAMES && g_rx[IDX_C].fin_seen && g_rx[IDX_C].seq_ok) + t->phase = P_DONE; + break; + case P_DONE: + t->result = 1; + uasync_stop(t->ua); + return; + } + + if (++t->ticks > MAX_TICKS) { + fprintf(stderr, "test: timeout in phase %d (B_rx=%d/%d fin=%d seqok=%d D_rx=%d/%d fin=%d seqok=%d C_rx=%d/%d fin=%d seqok=%d dupC=%u)\n", + (int)t->phase, + g_rx[IDX_B].frames, N_FRAMES, g_rx[IDX_B].fin_seen, g_rx[IDX_B].seq_ok, + g_rx[IDX_D].frames, N_FRAMES, g_rx[IDX_D].fin_seen, g_rx[IDX_D].seq_ok, + g_rx[IDX_C].frames, N_FRAMES, g_rx[IDX_C].fin_seen, g_rx[IDX_C].seq_ok, + radio_dup_dropped(C, gid)); + t->result = 2; + uasync_stop(t->ua); + return; + } + t->timer = uasync_set_timeout(t->ua, TICK_TB, t, mesh_tick, "mesh_tick"); +} + +/* ── main ── */ + +int main(int argc, char** argv) { + (void)argc; (void)argv; + debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + if (getenv("UTUN_TEST_DEBUG")) { + debug_set_category_level_by_name("radio", "debug"); + debug_set_category_level_by_name("bgp", "info"); + debug_set_category_level_by_name("member_sync", "info"); + debug_set_category_level_by_name("chat_sync", "info"); + } + utun_instance_set_tun_init_enabled(0); + srand((unsigned)time(NULL)); + + char dir[512]; snprintf(dir, sizeof(dir), "/tmp/utun_radio_mesh_XXXXXX"); + if (!mkdtemp(dir)) { fprintf(stderr, "mkdtemp failed\n"); return 1; } + setenv("UTUN_TEST_DIR", dir, 1); + + fill_shared(59000 + (getpid() % 2000)); + write_configs(dir); + + printf("=== test_radio_mesh ===\n"); fflush(stdout); + + TEST("single-thread: diamond (4 nodes) — дубль по 2 путям гасится кэшем (C получает ровно N)"); { + struct UASYNC* ua = uasync_create(); + if (!ua) { FAIL("uasync_create failed"); return 1; } + + struct mesh_ctx t; + memset(&t, 0, sizeof(t)); + t.ua = ua; + t.phase = P_WAIT_CONN; + + char cfg[4][512]; + const char* names[4] = { "a.conf", "b.conf", "d.conf", "c.conf" }; + for (int i = 0; i < 4; i++) snprintf(cfg[i], sizeof(cfg[i]), "%s/%s", dir, names[i]); + + for (int i = 0; i < 4; i++) { + t.inst[i] = utun_instance_create(ua, cfg[i]); + if (!t.inst[i]) { FAIL("instance create failed i=%d", i); uasync_destroy(ua, 0); return 1; } + } + for (int i = 0; i < 4; i++) utun_instance_init(t.inst[i]); + + for (int i = 0; i < 4; i++) radio_set_frame_cb(t.inst[i], rx_frame_cb, &g_rx[i]); + + t.timer = uasync_set_timeout(ua, TICK_TB, &t, mesh_tick, "mesh_tick"); + uasync_mainloop(ua); + + int ok = (t.result == 1); + uint64_t gid = strtoull(CH_ID, NULL, 10); + uint32_t dup_c = radio_dup_dropped(t.inst[IDX_C], gid); + + if (ok && dup_c == 0) { + /* кадры не задвоились, но и дедупликация не сработала — топология не дала дублей */ + printf("WARN: dupC=0 — diamond не породил дублей, кэш не проверен\n"); + } + + for (int i = 0; i < 4; i++) if (t.inst[i]) utun_instance_destroy(t.inst[i]); + uasync_destroy(ua, 0); + + if (ok && dup_c > 0) OK(); + else if (ok) FAIL("result ok, но дедупликация не сработала (dupC=%u)", dup_c); + else FAIL("result=%d phase=%d dupC=%u", t.result, (int)t.phase, dup_c); + } + + char cmd[512]; snprintf(cmd, sizeof(cmd), "rm -rf %s", dir); (void)!system(cmd); + + printf("\n%d/%d passed, %d failed\n", G_PASSED, G_TOTAL, G_FAILED); + return G_FAILED > 0 ? 1 : 0; +}