From a1ecab997ebe6367810b22e3d388cd083a14f6e9 Mon Sep 17 00:00:00 2001 From: Evgeny Date: Sun, 26 Jul 2026 17:42:27 +0300 Subject: [PATCH] merkle_sync: fix 4 protocol bugs + add test suite (Stages 1-4) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - g_merkle singleton → per-instance msync (utun_instance.h) - level_prefix: off-by-one shift (63→64) - _handle_hashes: differs |= for shared buckets with different hashes - _handle_batch: sub_pl double-level-offset (next_lvl→lvl) - pending[] array replaces pending_requests counter - member_sync: mask_shift 63→64 alignment - test_merkle_sync.c: prefix arithmetic, tree building, 2-instance integration, randomized 50-iter sync --- src/utun_instance.h | 2 + tests/Makefile.am | 5 + tests/test_merkle_sync.c | 982 ++++++++++++++++++++++++++ tools/chatgui/transport/member_sync.c | 8 +- tools/chatgui/transport/merkle_sync.c | 123 ++-- 5 files changed, 1071 insertions(+), 49 deletions(-) create mode 100644 tests/test_merkle_sync.c diff --git a/src/utun_instance.h b/src/utun_instance.h index 479b2553..8e9f8583 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -46,6 +46,7 @@ struct CONN_MGR; struct ETCP_CONNECT; struct DB_SYNC; struct NAT_DETECTION; +struct merkle_sync; struct NETWORK_ENTRY { uint64_t id; // 56-bit (offset 0 = index key) @@ -78,6 +79,7 @@ struct UTUN_INSTANCE { struct TOPO_GROUPS* topo_groups; // Groups module for topology exchange sqlite3* topo_sqlite_db; // Shared SQLite DB (nodes/channels/peers) struct NAT_DETECTION* nat_det; // NAT detection module + struct merkle_sync* msync; // Merkle tree sync module (per-instance) // Identification uint64_t node_id; diff --git a/tests/Makefile.am b/tests/Makefile.am index 63628172..b7006df0 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -50,6 +50,7 @@ check_PROGRAMS = \ test_conn_mgr \ test_etcp_connect \ test_db_sync \ + test_merkle_sync \ test_chat_sync_stress \ test_stcp_traffic \ test_bbr_integration \ @@ -281,6 +282,10 @@ test_db_sync_SOURCES = test_db_sync.c test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_merkle_sync_SOURCES = test_merkle_sync.c $(top_srcdir)/tools/chatgui/transport/merkle_sync.c +test_merkle_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tools/chatgui/transport +test_merkle_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_chat_sync_stress_SOURCES = test_chat_sync_stress.c test_chat_sync_stress_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_chat_sync_stress_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c new file mode 100644 index 00000000..117d18ae --- /dev/null +++ b/tests/test_merkle_sync.c @@ -0,0 +1,982 @@ +#include +#include +#include +#include +#include +#include +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/platform_compat.h" +#include "../lib/socket_compat.h" +#include "merkle_sync.h" +#include "../src/utun_instance.h" +#include "../src/config_updater.h" +#include "../src/tun_if.h" +#include "etcp.h" +#include "etcp_connections.h" +#include "test_utils.h" +#include +#include + +/* merkle_sync private struct — needed for Stage 2 unit tests */ +struct ms_session { + struct ms_session* next; + char ns[64]; + uint64_t peer; + uint8_t active; + uint8_t synced; + void* done_cb; /* merkle_sync_done_cb */ + void* cb_arg; +}; + +/* merkle_sync private struct — needed for Stage 2 unit tests. + * Must match struct merkle_sync in merkle_sync.c exactly. */ +struct merkle_sync { + struct UTUN_INSTANCE* inst; + uint8_t svc_id; + const struct merkle_sync_data_ops* ops; + void* data_ctx; + struct ms_session* sessions; + uint8_t initialized; + void* bg_timer; +}; + +static int tests_run = 0; +static int tests_failed = 0; +static int stage_failures = 0; + +#define TEST(name) do { tests_run++; printf(" %s: ", name); } while(0) +#define PASS() do { printf("PASS\n"); } while(0) +#define FAIL(fmt, ...) do { printf("FAIL — " fmt "\n", ##__VA_ARGS__); tests_failed++; stage_failures++; } while(0) + +/* ── Stage 2: mock data store ── */ + +#define MS_MAX_ITEMS 4096 + +struct ms_item { uint64_t key; uint32_t val; }; + +struct ms_data { + struct ms_item items[MS_MAX_ITEMS]; + int count; /* sorted by key */ +}; + +static int _cmp_item(const void* a, const void* b) { + uint64_t ka = ((const struct ms_item*)a)->key, kb = ((const struct ms_item*)b)->key; + return ka < kb ? -1 : ka > kb ? 1 : 0; +} + +static void _data_sort(struct ms_data* d) { qsort(d->items, (size_t)d->count, sizeof(struct ms_item), _cmp_item); } + +static void _data_init(struct ms_data* d) { memset(d, 0, sizeof(*d)); } + +static void _data_insert(struct ms_data* d, uint64_t key, uint32_t val) { + for (int i = 0; i < d->count; i++) + if (d->items[i].key == key) { d->items[i].val = val; return; } + if (d->count >= MS_MAX_ITEMS) return; + d->items[d->count].key = key; d->items[d->count].val = val; + d->count++; +} + +static int _data_count_in_prefix(struct ms_data* d, uint8_t level, uint64_t prefix64) { + int mask_shift = 64 - (int)level * 5; + uint64_t mask = (mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX; + int cnt = 0; + for (int i = 0; i < d->count; i++) + if ((d->items[i].key & mask) == prefix64) cnt++; + return cnt; +} + +/* ── Stage 2: mock merkle_sync_data_ops ── */ + +struct ms_test_ctx { + struct ms_data* data; + struct UTUN_INSTANCE* inst; + char tag; /* 'A'/'B'/'C' for debug */ +}; + +static void _compute_item_hash(uint64_t key, uint32_t val, uint8_t out[MT_HASH_SIZE]) { + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); + EVP_DigestUpdate(ctx, &key, 8); + EVP_DigestUpdate(ctx, &val, 4); + EVP_DigestFinal_ex(ctx, out, NULL); + EVP_MD_CTX_free(ctx); +} + +static int _bucket_hash(void* ctx, const char* ns, uint8_t level, uint64_t prefix64, EVP_MD_CTX* sha_ctx) { + (void)ns; + struct ms_data* d = ((struct ms_test_ctx*)ctx)->data; + int mask_shift = 64 - (int)level * 5; + uint64_t mask = (level == 0) ? 0 : ((mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX); + int count = 0; + for (int i = 0; i < d->count; i++) { + if ((d->items[i].key & mask) != prefix64) continue; + uint8_t ih[MT_HASH_SIZE]; _compute_item_hash(d->items[i].key, d->items[i].val, ih); + EVP_DigestUpdate(sha_ctx, ih, MT_HASH_SIZE); + count++; + } + return count; +} + +static int _get_items(void* ctx, const char* ns, uint8_t level, uint64_t prefix, + uint8_t pbytes, uint8_t* buf, size_t* len) { + (void)ns; (void)pbytes; + struct ms_data* d = ((struct ms_test_ctx*)ctx)->data; + int mask_shift = 64 - (int)level * 5; + uint64_t mask = (level == 0) ? 0 : ((mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX); + uint16_t count = 0; size_t off = 2; + for (int i = 0; i < d->count; i++) { + if ((d->items[i].key & mask) != prefix) continue; + if (off + 12 > *len) return -2; + memcpy(buf + off, &d->items[i].key, 8); off += 8; + memcpy(buf + off, &d->items[i].val, 4); off += 4; + count++; + } + memcpy(buf, &count, 2); + *len = off; + return 0; +} + +static int _apply_items(void* ctx, const char* ns, const uint8_t* data, size_t len) { + struct ms_test_ctx* tc = (struct ms_test_ctx*)ctx; + struct ms_data* d = tc->data; + if (len < 2) return 0; + uint16_t count; memcpy(&count, data, 2); + const uint8_t* p = data + 2; + int changed = 0; + for (uint16_t i = 0; i < count; i++) { + uint64_t key; memcpy(&key, p, 8); p += 8; + uint32_t val; memcpy(&val, p, 4); p += 4; + _data_insert(d, key, val); + changed++; + } + if (changed && tc->inst) { + _data_sort(d); + for (uint16_t i = 0; i < count; i++) { + uint64_t k; memcpy(&k, data + 2 + (size_t)i * 12, 8); + merkle_sync_recompute_path(tc->inst, ns, k); + } + } + return 0; +} + +static int _apply_update(void* ctx, const char* ns, uint64_t key, uint8_t type, + const uint8_t* data, size_t len) { + (void)ctx; (void)ns; (void)key; (void)type; (void)data; (void)len; + return 0; +} + +static const struct merkle_sync_data_ops g_test_ops = { + .update_bucket_hash = _bucket_hash, + .get_items = _get_items, + .apply_items = _apply_items, + .apply_update = _apply_update, +}; + +/* ── Stage 1: prefix arithmetic ── */ + +static void test_prefix_bytes(void) { + TEST("prefix_bytes(1)==1"); + if (merkle_sync_prefix_bytes(1) == 1) PASS(); else FAIL("got %d", merkle_sync_prefix_bytes(1)); + + TEST("prefix_bytes(2)==2"); + if (merkle_sync_prefix_bytes(2) == 2) PASS(); else FAIL("got %d", merkle_sync_prefix_bytes(2)); + + TEST("prefix_bytes(3)==2"); + if (merkle_sync_prefix_bytes(3) == 2) PASS(); else FAIL("got %d", merkle_sync_prefix_bytes(3)); + + TEST("prefix_bytes(4)==3"); + if (merkle_sync_prefix_bytes(4) == 3) PASS(); else FAIL("got %d", merkle_sync_prefix_bytes(4)); + + TEST("prefix_bytes(5)==4"); + if (merkle_sync_prefix_bytes(5) == 4) PASS(); else FAIL("got %d", merkle_sync_prefix_bytes(5)); +} + +static void test_level_prefix_basic(void) { + uint64_t key = 0xFEDCBA9876543210ULL; + + TEST("level_prefix(0xFEDCBA.., 1)==0xF800.."); + uint64_t got = merkle_sync_level_prefix(key, 1); + if (got == 0xF800000000000000ULL) PASS(); else FAIL("expected f800000000000000 got %016llx", (unsigned long long)got); + + TEST("level_prefix(0,3)==0"); + got = merkle_sync_level_prefix(0, 3); + if (got == 0) PASS(); else FAIL("got %016llx", (unsigned long long)got); + + TEST("level_prefix(MAX,1)==0xF800000000000000"); + got = merkle_sync_level_prefix(UINT64_MAX, 1); + if (got == 0xF800000000000000ULL) PASS(); else FAIL("got %016llx", (unsigned long long)got); + + TEST("level_prefix(MAX,5)==0xFFFFFF8000000000"); + got = merkle_sync_level_prefix(UINT64_MAX, 5); + if (got == 0xFFFFFF8000000000ULL) PASS(); else FAIL("got %016llx", (unsigned long long)got); + + TEST("level_prefix(0,0)==0"); + got = merkle_sync_level_prefix(0, 0); + if (got == 0) PASS(); else FAIL("got %016llx", (unsigned long long)got); +} + +static void test_level_prefix_idempotent(void) { + TEST("level_prefix idempotent L3"); + uint64_t key = 0x0123456789ABCDEFULL; + uint64_t p = merkle_sync_level_prefix(key, 3); + uint64_t p2 = merkle_sync_level_prefix(p, 3); + if (p == p2) PASS(); else FAIL("p=%016llx p2=%016llx", (unsigned long long)p, (unsigned long long)p2); +} + +static void test_level_prefix_bucket_partition(void) { + TEST("same bucket → same prefix L1..L2"); + uint64_t k1 = 0x8900000000000000ULL; + uint64_t k2 = 0x8911111111111111ULL; + int ok = 1; + for (uint8_t l = 1; l <= 2; l++) + if (merkle_sync_level_prefix(k1, l) != merkle_sync_level_prefix(k2, l)) ok = 0; + if (merkle_sync_level_prefix(k1, 3) == merkle_sync_level_prefix(k2, 3)) ok = 0; + if (ok) PASS(); else FAIL("partition broken"); + + TEST("different L1 bucket → different prefix L1"); + uint64_t a = 0x0000000000000000ULL; + uint64_t b = 0x0800000000000000ULL; + if (merkle_sync_level_prefix(a, 1) != merkle_sync_level_prefix(b, 1)) PASS(); else FAIL("should differ"); +} + +/* ── Stage 2: tree building tests ── */ + +static void _ensure_table(sqlite3* db) { + sqlite3_exec(db, + "CREATE TABLE IF NOT EXISTS merkle_tree_hash (" + " namespace TEXT NOT NULL," + " level INTEGER NOT NULL CHECK(level BETWEEN 1 AND 5)," + " prefix64 INTEGER NOT NULL," + " hash BLOB NOT NULL," + " member_count INTEGER NOT NULL," + " PRIMARY KEY (namespace, level, prefix64))", + NULL, NULL, NULL); +} + +static int _db_count_rows(sqlite3* db, const char* ns) { + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(db, "SELECT COUNT(*) FROM merkle_tree_hash WHERE namespace=?", -1, &st, NULL); + sqlite3_bind_text(st, 1, ns, -1, SQLITE_STATIC); + int n = 0; + if (sqlite3_step(st) == SQLITE_ROW) n = sqlite3_column_int(st, 0); + sqlite3_finalize(st); + return n; +} + +static int _db_compare_trees(sqlite3* da, sqlite3* db, const char* ns) { + sqlite3_stmt* sa = NULL, *sb = NULL; + sqlite3_prepare_v2(da, "SELECT level,prefix64,hash,member_count FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", + -1, &sa, NULL); + sqlite3_prepare_v2(db, "SELECT level,prefix64,hash,member_count FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", + -1, &sb, NULL); + sqlite3_bind_text(sa, 1, ns, -1, SQLITE_STATIC); + sqlite3_bind_text(sb, 1, ns, -1, SQLITE_STATIC); + int mismatch = 0; + while (1) { + int ra = sqlite3_step(sa), rb = sqlite3_step(sb); + if (ra == SQLITE_DONE && rb == SQLITE_DONE) break; + if (ra == SQLITE_DONE) { mismatch++; break; } + if (rb == SQLITE_DONE) { mismatch++; break; } + if (sqlite3_column_int(sa,0) != sqlite3_column_int(sb,0) || + (uint64_t)sqlite3_column_int64(sa,1) != (uint64_t)sqlite3_column_int64(sb,1)) { + mismatch++; break; + } + if (sqlite3_column_int(sa,3) != sqlite3_column_int(sb,3)) { mismatch++; break; } + if (memcmp(sqlite3_column_blob(sa,2), sqlite3_column_blob(sb,2), MT_HASH_SIZE) != 0) { mismatch++; break; } + } + sqlite3_finalize(sa); sqlite3_finalize(sb); + return mismatch; +} + +static void test_tree_empty_ns(void) { + TEST("empty namespace → no rows in merkle_tree_hash"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + merkle_sync_recompute_path(&inst, "test", 0x1234ULL); + if (_db_count_rows(db, "test") == 0) PASS(); else FAIL("got %d rows", _db_count_rows(db, "test")); + + sqlite3_close(db); +} + +static void test_tree_single_item(void) { + TEST("single item → 5 rows (L1..L5) with non-zero hash"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + uint64_t key = 0xABCD000000000000ULL; + _data_insert(&data, key, 42); + _data_sort(&data); + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + merkle_sync_recompute_path(&inst, "test", key); + int rows = _db_count_rows(db, "test"); + if (rows != 5) { FAIL("expected 5 rows got %d", rows); sqlite3_close(db); return; } + + /* check hashes are non-zero */ + int ok = 1; + for (uint8_t lv = 1; lv <= 5 && ok; lv++) { + uint64_t pf = merkle_sync_level_prefix(key, lv); + const uint8_t* h = merkle_sync_get_hash(&inst, "test", lv, pf); + int zero = 1; + for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0; + if (zero) { FAIL("L%d zero hash", lv); ok = 0; } + } + if (ok) PASS(); + + sqlite3_close(db); +} + +static void test_tree_multi_item_same_bucket(void) { + TEST("two items same L1 bucket → 5 rows, L1 hash aggregates both"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + uint64_t k1 = 0xAA00000000000000ULL; + uint64_t k2 = 0xAB00000000000000ULL; /* bits 63-59: 10101 for both */ + _data_insert(&data, k1, 1); _data_insert(&data, k2, 2); + _data_sort(&data); + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + /* recompute for both */ + merkle_sync_recompute_path(&inst, "test", k1); + merkle_sync_recompute_path(&inst, "test", k2); + + uint64_t pf1 = merkle_sync_level_prefix(k1, 1); + uint64_t pf2 = merkle_sync_level_prefix(k2, 1); + if (pf1 != pf2) { FAIL("L1 prefixes should be equal"); sqlite3_close(db); return; } + + /* verify: both items contribute to same bucket hash */ + uint8_t expected_hash[MT_HASH_SIZE]; + { + EVP_MD_CTX* ctx = EVP_MD_CTX_new(); + EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); + uint8_t ih[MT_HASH_SIZE]; _compute_item_hash(k1, 1, ih); EVP_DigestUpdate(ctx, ih, MT_HASH_SIZE); + _compute_item_hash(k2, 2, ih); EVP_DigestUpdate(ctx, ih, MT_HASH_SIZE); + EVP_DigestFinal_ex(ctx, expected_hash, NULL); + EVP_MD_CTX_free(ctx); + } + const uint8_t* stored = merkle_sync_get_hash(&inst, "test", 1, pf1); + if (memcmp(expected_hash, stored, MT_HASH_SIZE) == 0) PASS(); else FAIL("hash mismatch"); + + sqlite3_close(db); +} + +static void test_tree_different_buckets(void) { + TEST("items in different L1 buckets → independent hashes"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + uint64_t kA = 0x0000000000000000ULL; /* bucket 0 */ + uint64_t kB = 0x0800000000000000ULL; /* bucket 1 */ + _data_insert(&data, kA, 100); _data_insert(&data, kB, 200); + _data_sort(&data); + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + merkle_sync_recompute_path(&inst, "ns", kA); + merkle_sync_recompute_path(&inst, "ns", kB); + + uint64_t pfA = merkle_sync_level_prefix(kA, 1); + uint64_t pfB = merkle_sync_level_prefix(kB, 1); + if (pfA == pfB) { FAIL("same L1 prefix"); sqlite3_close(db); return; } + + const uint8_t* hA = merkle_sync_get_hash(&inst, "ns", 1, pfA); + uint8_t copyA[MT_HASH_SIZE]; memcpy(copyA, hA, MT_HASH_SIZE); + const uint8_t* hB = merkle_sync_get_hash(&inst, "ns", 1, pfB); + uint8_t copyB[MT_HASH_SIZE]; memcpy(copyB, hB, MT_HASH_SIZE); + if (memcmp(copyA, copyB, MT_HASH_SIZE) != 0) PASS(); else FAIL("hashes should differ"); + + sqlite3_close(db); +} + +static void test_tree_delete_empty_bucket(void) { + TEST("delete last item → bucket deleted from DB, get_hash returns zero"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + uint64_t k = 0xCCCC000000000000ULL; + _data_insert(&data, k, 7); _data_sort(&data); + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + merkle_sync_recompute_path(&inst, "test", k); + if (_db_count_rows(db, "test") != 5) { FAIL("expected 5 rows"); sqlite3_close(db); return; } + + /* remove item and recompute */ + data.count = 0; + merkle_sync_recompute_path(&inst, "test", k); + if (_db_count_rows(db, "test") != 0) { FAIL("expected 0 rows got %d", _db_count_rows(db, "test")); sqlite3_close(db); return; } + + uint64_t pf = merkle_sync_level_prefix(k, 3); + const uint8_t* h = merkle_sync_get_hash(&inst, "test", 3, pf); + int zero = 1; + for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0; + if (zero) PASS(); else FAIL("hash not zero after delete"); + + sqlite3_close(db); +} + +static void test_tree_consistency(void) { + TEST("recompute vs stored hash consistency — random items"); + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + _ensure_table(db); + + static struct UTUN_INSTANCE inst; memset(&inst, 0, sizeof(inst)); + inst.topo_sqlite_db = db; + struct ms_data data; _data_init(&data); + + /* random items across the keyspace */ + srand(12345); + for (int i = 0; i < 50; i++) { + uint64_t k = ((uint64_t)rand() << 32) | (uint64_t)rand(); + _data_insert(&data, k, (uint32_t)rand()); + } + _data_sort(&data); + + struct ms_test_ctx ctx = { .data = &data, .inst = &inst }; + struct merkle_sync ms; memset(&ms, 0, sizeof(ms)); + ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; + inst.msync = &ms; + + for (int i = 0; i < data.count; i++) + merkle_sync_recompute_path(&inst, "ns", data.items[i].key); + + /* verify: recompute each bucket, compare with stored hash */ + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(db, + "SELECT level, prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level, prefix64", + -1, &st, NULL); + sqlite3_bind_text(st, 1, "ns", -1, SQLITE_STATIC); + int ok = 1; + while (sqlite3_step(st) == SQLITE_ROW && ok) { + uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); + uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); + + /* compute expected hash from data */ + EVP_MD_CTX* c = EVP_MD_CTX_new(); + EVP_DigestInit_ex(c, EVP_sha256(), NULL); + _bucket_hash(&ctx, "ns", lv, pf, c); + uint8_t expected[MT_HASH_SIZE]; + EVP_DigestFinal_ex(c, expected, NULL); + EVP_MD_CTX_free(c); + + const uint8_t* stored = merkle_sync_get_hash(&inst, "ns", lv, pf); + if (memcmp(expected, stored, MT_HASH_SIZE) != 0) { + FAIL("mismatch L%d P%016llx", lv, (unsigned long long)pf); + ok = 0; + } + } + sqlite3_finalize(st); + if (ok) PASS(); + + sqlite3_close(db); +} + +static int run_stage1(void) { + printf("--- Stage 1: prefix arithmetic ---\n"); + test_prefix_bytes(); + test_level_prefix_basic(); + test_level_prefix_idempotent(); + test_level_prefix_bucket_partition(); + return tests_failed; +} + +static int run_stage2(void) { + printf("--- Stage 2: tree building with in-memory SQLite ---\n"); + stage_failures = 0; + test_tree_empty_ns(); + test_tree_single_item(); + test_tree_multi_item_same_bucket(); + test_tree_different_buckets(); + test_tree_delete_empty_bucket(); + test_tree_consistency(); + return stage_failures; +} + +/* ── Stage 3-6: integration with UTUN_INSTANCE + ETCP ── */ + +#define INTG_TIMEOUT_TB 120000 +#define PHASE_TIMEOUT_TB 100000 +#define POLL_MS 1 +#define MS_SVC_ID 0x72 + +static char intg_temp_dir[] = "/tmp/utun_msync_XXXXXX"; +static struct UTUN_INSTANCE* i_a = NULL; +static struct UTUN_INSTANCE* i_b = NULL; +static struct UTUN_INSTANCE* i_c = NULL; +static struct UASYNC* i_ua = NULL; +static int i_phase = 0; +static void* i_timeout_id = NULL; + +static struct ms_data data_a, data_b, data_c; +static struct ms_test_ctx ctx_a, ctx_b, ctx_c; + +static int done_sync = 0; +static int _cond_done(void) { return done_sync; } + +static void _on_sync_done(uint64_t peer, const char* ns, int result, void* arg) { + (void)arg; done_sync = 1; + printf(" sync_done: peer=%016llx ns=%s result=%d\n", (unsigned long long)peer, ns, result); +} + +static int _write_cfg(const char* path, const char* fmt, ...) { + va_list ap; FILE* f = fopen(path, "w"); + if (!f) return -1; + va_start(ap, fmt); vfprintf(f, fmt, ap); va_end(ap); + fclose(f); return 0; +} + +static char* _get_pubkey(const char* path) { + struct utun_config* cfg = parse_config(path); + if (!cfg) return NULL; + char* pub = u_strdup(cfg->global.my_public_key_hex); + free_config(cfg); return pub; +} + +static void _intg_timeout(void* arg) { (void)arg; printf(" GLOBAL TIMEOUT\n"); i_phase = 2; } + +static int _cond_links_up(void) { + if (!i_a || !i_b) return 0; + int links = 0; + struct ll_entry* e = i_a->connections->head; + while (e) { + struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; + struct ETCP_LINK* l = ce->conn->links; + while (l) { if (l->initialized) links++; l = l->next; } + e = e->next; + } + return links >= (i_c ? 2 : 1); +} + +static int _wait_for(const char* desc, int (*cond)(void), int timeout_tb) { + uint64_t start = get_time_tb(); + while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && i_phase == 0) + uasync_poll(i_ua, POLL_MS); + if (cond()) return 1; + if (i_phase == 0) { printf(" TIMEOUT: %s\n", desc); i_phase = 2; } + return 0; +} + +static void _intg_setup_contexts(void) { + memset(&data_a, 0, sizeof(data_a)); memset(&data_b, 0, sizeof(data_b)); memset(&data_c, 0, sizeof(data_c)); + ctx_a.data = &data_a; ctx_a.inst = NULL; ctx_a.tag = 'A'; + ctx_b.data = &data_b; ctx_b.inst = NULL; ctx_b.tag = 'B'; + ctx_c.data = &data_c; ctx_c.inst = NULL; ctx_c.tag = 'C'; + done_sync = 0; +} + +static int _intg_create_cfgs(int* port_a, int* port_b, int* port_c, + char* cfg_a, char* cfg_b, char* cfg_c, size_t sz) { + /* reset template (mkdtemp modifies it) */ + strcpy(intg_temp_dir, "/tmp/utun_msync_XXXXXX"); + if (test_mkdtemp(intg_temp_dir) != 0) { printf(" mkdtemp failed\n"); return -1; } + int base = 43000 + (getpid() % 15000); + *port_a = base; *port_b = base + 1; *port_c = base + 2; + snprintf(cfg_a, sz, "%s/a.conf", intg_temp_dir); + snprintf(cfg_b, sz, "%s/b.conf", intg_temp_dir); + snprintf(cfg_c, sz, "%s/c.conf", intg_temp_dir); + char dbp[256]; + snprintf(dbp, sizeof(dbp), "%s/db_a", intg_temp_dir); utun_mkdir(dbp, 0755); + snprintf(dbp, sizeof(dbp), "%s/db_b", intg_temp_dir); utun_mkdir(dbp, 0755); + snprintf(dbp, sizeof(dbp), "%s/db_c", intg_temp_dir); utun_mkdir(dbp, 0755); + + _write_cfg(cfg_a, + "[global]\nmy_node_id=0xAAAAAAAAAAAAAAAA\ntun_ip=10.220.0.1/24\n" + "tun_ifname=tun220\nkeepalive_adaptive=0\ndb_path=%s/db_a\n\n" + "[server: srv_a]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[allowed_keys]\nallow_all=1\n", intg_temp_dir, *port_a); + config_ensure_keys_and_node_id(cfg_a); + char* pub_a = _get_pubkey(cfg_a); if (!pub_a) return -1; + + _write_cfg(cfg_b, + "[global]\nmy_node_id=0xBBBBBBBBBBBBBBBB\ntun_ip=10.220.0.2/24\n" + "tun_ifname=tun221\nkeepalive_adaptive=0\ndb_path=%s/db_b\n\n" + "[server: srv_b]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_b:127.0.0.1:%d\n\n" + "[allowed_keys]\nallow_all=1\n", intg_temp_dir, *port_b, pub_a, *port_a); + config_ensure_keys_and_node_id(cfg_b); + + _write_cfg(cfg_c, + "[global]\nmy_node_id=0xCCCCCCCCCCCCCCCC\ntun_ip=10.220.0.3/24\n" + "tun_ifname=tun222\nkeepalive_adaptive=0\ndb_path=%s/db_c\n\n" + "[server: srv_c]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_c:127.0.0.1:%d\n\n" + "[allowed_keys]\nallow_all=1\n", intg_temp_dir, *port_c, pub_a, *port_a); + config_ensure_keys_and_node_id(cfg_c); + + u_free(pub_a); + return 0; +} + +static void _intg_cleanup(void) { + if (i_timeout_id) { uasync_cancel_timeout(i_ua, i_timeout_id); i_timeout_id = NULL; } + if (i_c) { merkle_sync_destroy(i_c); i_c->running = 0; utun_instance_destroy(i_c); i_c = NULL; } + if (i_b) { merkle_sync_destroy(i_b); i_b->running = 0; utun_instance_destroy(i_b); i_b = NULL; } + if (i_a) { merkle_sync_destroy(i_a); i_a->running = 0; utun_instance_destroy(i_a); i_a = NULL; } + if (i_ua) { uasync_destroy(i_ua, 0); i_ua = NULL; } + char pa[320]; + snprintf(pa, sizeof(pa), "%s/a.conf", intg_temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/b.conf", intg_temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/c.conf", intg_temp_dir); unlink(pa); + for (const char* d = "abc"; *d; d++) { + snprintf(pa, sizeof(pa), "%s/db_%c/chats.db", intg_temp_dir, *d); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_%c/chats.db-wal", intg_temp_dir, *d); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_%c/chats.db-shm", intg_temp_dir, *d); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_%c", intg_temp_dir, *d); test_rmdir(pa); + } + test_rmdir(intg_temp_dir); +} + +static int _intg_init_two(void) { + i_phase = 0; i_timeout_id = NULL; i_a = i_b = i_c = NULL; + int port_a, port_b, port_c; + char cfg_a[256], cfg_b[256], cfg_c[256]; + if (_intg_create_cfgs(&port_a, &port_b, &port_c, cfg_a, cfg_b, cfg_c, sizeof(cfg_a)) != 0) return -1; + utun_instance_set_tun_init_enabled(0); + i_ua = uasync_create(); if (!i_ua) { _intg_cleanup(); return -1; } + i_a = utun_instance_create(i_ua, cfg_a); + i_b = utun_instance_create(i_ua, cfg_b); + if (!i_a || !i_b || utun_instance_init(i_a) != 0 || utun_instance_init(i_b) != 0) + { _intg_cleanup(); return -1; } + i_timeout_id = uasync_set_timeout(i_ua, INTG_TIMEOUT_TB, NULL, _intg_timeout, "ms_intg_timeout"); + if (!_wait_for("links up", _cond_links_up, PHASE_TIMEOUT_TB)) { _intg_cleanup(); return -1; } + + _intg_setup_contexts(); + ctx_a.inst = i_a; ctx_b.inst = i_b; + if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 || + merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0) + { _intg_cleanup(); return -1; } + printf(" 2 instances ready\n"); + return 0; +} + +static void _intg_ins_many(struct ms_data* d, struct UTUN_INSTANCE* inst, int base, int count) { + for (int i = 0; i < count; i++) { + uint64_t key = (uint64_t)(base + i) * 0x100000000000000ULL; + _data_insert(d, key, (uint32_t)(base + i)); + } + _data_sort(d); + for (int i = 0; i < d->count; i++) + merkle_sync_recompute_path(inst, "test", d->items[i].key); +} + +static int _intg_compare_data(struct ms_data* a, struct ms_data* b) { + if (a->count != b->count) { printf(" count mismatch: %d vs %d\n", a->count, b->count); return -1; } + for (int i = 0; i < a->count; i++) { + if (a->items[i].key != b->items[i].key || a->items[i].val != b->items[i].val) { + printf(" data[%d] mismatch\n", i); return -1; + } + } + return 0; +} + +static int _intg_verify_all(struct UTUN_INSTANCE* a, struct UTUN_INSTANCE* b, + struct ms_data* da, struct ms_data* db, const char* ns) { + if (_db_compare_trees(a->topo_sqlite_db, b->topo_sqlite_db, ns) != 0) { printf(" TREE MISMATCH\n"); return -1; } + if (_intg_compare_data(da, db) != 0) { printf(" DATA MISMATCH\n"); return -1; } + return 0; +} + +/* ── Stage 3: integration tests ── */ + +static void test_peer_empty(void) { + TEST("peer empty — B has 5 items, A empty, sync A->B"); + if (_intg_init_two() != 0) { FAIL("setup failed"); return; } + _intg_ins_many(&data_b, i_b, 0, 5); + done_sync = 0; + merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); + if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB)) + { FAIL("sync timeout"); _intg_cleanup(); return; } + printf(" A data=%d B data=%d, A tree=%d B tree=%d\n", + data_a.count, data_b.count, + _db_count_rows(i_a->topo_sqlite_db, "test"), + _db_count_rows(i_b->topo_sqlite_db, "test")); + if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) + { FAIL("verify failed"); _intg_cleanup(); return; } + PASS(); + _intg_cleanup(); +} + +static void test_hash_match(void) { + TEST("hash match — identical data, sync immediate"); + if (_intg_init_two() != 0) { FAIL("setup failed"); return; } + _intg_ins_many(&data_a, i_a, 0, 5); + _intg_ins_many(&data_b, i_b, 0, 5); + done_sync = 0; + merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); + for (int i = 0; i < 200 && !done_sync && i_phase == 0; i++) uasync_poll(i_ua, 1); + if (!done_sync) { FAIL("sync did not complete quickly"); _intg_cleanup(); return; } + if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) + { FAIL("verify failed"); _intg_cleanup(); return; } + PASS(); + _intg_cleanup(); +} + +static void test_divergence_merge(void) { + TEST("divergence merge — A has 3 items, B has 2 different"); + if (_intg_init_two() != 0) { FAIL("setup failed"); return; } + _intg_ins_many(&data_a, i_a, 0, 3); + _intg_ins_many(&data_b, i_b, 10, 2); + done_sync = 0; + merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); + if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB)) + { FAIL("sync timeout"); _intg_cleanup(); return; } + if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) + { FAIL("verify failed"); _intg_cleanup(); return; } + PASS(); + _intg_cleanup(); +} + +static int run_stage3(void) { + printf("--- Stage 3: integration 2 instances, basic scenarios ---\n"); + stage_failures = 0; + test_peer_empty(); + test_hash_match(); + test_divergence_merge(); + return stage_failures; +} + +/* ── Stage 4: randomized 2-instance sync ── */ + +static void test_randomized_two(void) { + TEST("randomized 2-instance — 50 iterations"); + srand(42); + int iter; + for (iter = 0; iter < 50 && stage_failures == 0; iter++) { + if (_intg_init_two() != 0) { FAIL("setup iter %d", iter); break; } + int na = rand() % 101, nb = rand() % 101; + for (int i = 0; i < na; i++) { + uint64_t k = ((uint64_t)rand() << 32) | (uint64_t)rand(); + _data_insert(&data_a, k, (uint32_t)rand()); + } + for (int i = 0; i < nb; i++) { + uint64_t k = ((uint64_t)rand() << 32) | (uint64_t)rand(); + _data_insert(&data_b, k, (uint32_t)rand()); + } + _data_sort(&data_a); _data_sort(&data_b); + for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, "rnd", data_a.items[i].key); + for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.items[i].key); + + done_sync = 0; + merkle_sync_start(i_a, i_b->node_id, "rnd", _on_sync_done, NULL); + if (!_wait_for("sync", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout iter %d", iter); _intg_cleanup(); break; } + if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "rnd") != 0) { + printf(" A data=%d B data=%d A tree=%d B tree=%d\n", + data_a.count, data_b.count, + _db_count_rows(i_a->topo_sqlite_db, "rnd"), + _db_count_rows(i_b->topo_sqlite_db, "rnd")); + /* dump items only in A */ + for (int ai = 0; ai < data_a.count; ai++) { + int found = 0; + for (int bi = 0; bi < data_b.count; bi++) + if (data_a.items[ai].key == data_b.items[bi].key) { found = 1; break; } + if (!found) printf(" only A: key=%016llx val=%u\n", (unsigned long long)data_a.items[ai].key, data_a.items[ai].val); + } + for (int bi = 0; bi < data_b.count; bi++) { + int found = 0; + for (int ai = 0; ai < data_a.count; ai++) + if (data_b.items[bi].key == data_a.items[ai].key) { found = 1; break; } + if (!found) printf(" only B: key=%016llx val=%u\n", (unsigned long long)data_b.items[bi].key, data_b.items[bi].val); + } + FAIL("verify iter %d", iter); _intg_cleanup(); break; + } + + /* consistency check: recompute each bucket hash, verify stored matches */ + int ok = 1; + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(i_a->topo_sqlite_db, + "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", + -1, &st, NULL); + sqlite3_bind_text(st, 1, "rnd", -1, SQLITE_STATIC); + while (sqlite3_step(st) == SQLITE_ROW && ok) { + uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); + uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); + EVP_MD_CTX* c = EVP_MD_CTX_new(); + EVP_DigestInit_ex(c, EVP_sha256(), NULL); + _bucket_hash(&ctx_a, "rnd", lv, pf, c); + uint8_t expected[MT_HASH_SIZE]; + EVP_DigestFinal_ex(c, expected, NULL); EVP_MD_CTX_free(c); + const uint8_t* stored = merkle_sync_get_hash(i_a, "rnd", lv, pf); + uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE); + if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { ok = 0; } + } + sqlite3_finalize(st); + if (!ok) { FAIL("tree consistency iter %d", iter); _intg_cleanup(); break; } + _intg_cleanup(); + } + if (iter == 50) PASS(); +} + +/* ── Stage 5: randomized 3-instance star-topology sync ── */ + +static int _intg_init_three(void) { + i_phase = 0; i_timeout_id = NULL; i_a = i_b = i_c = NULL; + int port_a, port_b, port_c; + char cfg_a[256], cfg_b[256], cfg_c[256]; + if (_intg_create_cfgs(&port_a, &port_b, &port_c, cfg_a, cfg_b, cfg_c, sizeof(cfg_a)) != 0) return -1; + utun_instance_set_tun_init_enabled(0); + i_ua = uasync_create(); if (!i_ua) { _intg_cleanup(); return -1; } + i_a = utun_instance_create(i_ua, cfg_a); + i_b = utun_instance_create(i_ua, cfg_b); + i_c = utun_instance_create(i_ua, cfg_c); + if (!i_a || !i_b || !i_c || utun_instance_init(i_a) != 0 || utun_instance_init(i_b) != 0 || utun_instance_init(i_c) != 0) + { _intg_cleanup(); return -1; } + i_timeout_id = uasync_set_timeout(i_ua, INTG_TIMEOUT_TB, NULL, _intg_timeout, "ms_intg_timeout"); + if (!_wait_for("links up", _cond_links_up, PHASE_TIMEOUT_TB)) { _intg_cleanup(); return -1; } + _intg_setup_contexts(); + ctx_a.inst = i_a; ctx_b.inst = i_b; ctx_c.inst = i_c; + if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 || + merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0 || + merkle_sync_init(i_c, MS_SVC_ID, &g_test_ops, &ctx_c) != 0) + { _intg_cleanup(); return -1; } + printf(" 3 instances ready\n"); + return 0; +} + +static int _cond_all_synced(void) { + if (data_a.count != data_b.count || data_b.count != data_c.count) return 0; + if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0) return 0; + if (_db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, "rnd") != 0) return 0; + return 1; /* all three match — synced */ +} + +static void test_randomized_three(void) { + TEST("randomized 3-instance star — 30 iterations"); + srand(1234); + int iter; + for (iter = 0; iter < 30 && stage_failures == 0; iter++) { + if (_intg_init_three() != 0) { FAIL("setup iter %d", iter); break; } + int na = rand() % 51, nb = rand() % 51, nc = rand() % 51; + for (int i = 0; i < na; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_a, k, (uint32_t)rand()); } + for (int i = 0; i < nb; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_b, k, (uint32_t)rand()); } + for (int i = 0; i < nc; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_c, k, (uint32_t)rand()); } + _data_sort(&data_a); for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, "rnd", data_a.items[i].key); + _data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.items[i].key); + _data_sort(&data_c); for (int i = 0; i < data_c.count; i++) merkle_sync_recompute_path(i_c, "rnd", data_c.items[i].key); + + done_sync = 0; + merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); + merkle_sync_start(i_c, i_a->node_id, "rnd", _on_sync_done, NULL); + if (!_wait_for("sync", _cond_all_synced, PHASE_TIMEOUT_TB)) { FAIL("timeout iter %d", iter); _intg_cleanup(); break; } + if (_intg_compare_data(&data_a, &data_b) != 0 || _intg_compare_data(&data_b, &data_c) != 0) + { FAIL("data mismatch iter %d", iter); _intg_cleanup(); break; } + if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0 || + _db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, "rnd") != 0) + { FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; } + _intg_cleanup(); + } + if (iter == 30) PASS(); +} + +/* ── Stage 6: protocol edge cases ── */ + +static void test_cancel(void) { + TEST("cancel — start sync, immediately cancel, done_cb NOT called"); + if (_intg_init_two() != 0) { FAIL("setup failed"); return; } + _intg_ins_many(&data_b, i_b, 0, 5); + done_sync = 0; + merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); + uasync_poll(i_ua, 1); + merkle_sync_cancel(i_a, i_b->node_id, "test"); + for (int i = 0; i < 100 && i_phase == 0; i++) uasync_poll(i_ua, 10); + if (!done_sync) PASS(); else FAIL("done_cb was called"); + _intg_cleanup(); +} + +static int sw_cb1_called = 0, sw_cb2_called = 0; +static void _sw_done1(uint64_t p, const char* n, int r, void* a) { (void)p;(void)n;(void)r;(void)a; sw_cb1_called = 1; } +static void _sw_done2(uint64_t p, const char* n, int r, void* a) { (void)p;(void)n;(void)r;(void)a; sw_cb2_called = 1; } + +static void test_start_overwrite(void) { + TEST("start overwrite — cb1 replaced by cb2, only cb2 fires"); + if (_intg_init_two() != 0) { FAIL("setup failed"); return; } + _intg_ins_many(&data_b, i_b, 0, 3); + sw_cb1_called = sw_cb2_called = 0; + merkle_sync_start(i_a, i_b->node_id, "test", _sw_done1, NULL); + merkle_sync_start(i_a, i_b->node_id, "test", _sw_done2, NULL); + for (int i = 0; i < 500 && !sw_cb2_called && i_phase == 0; i++) uasync_poll(i_ua, 10); + if (sw_cb2_called) PASS(); else FAIL("cb1=%d cb2=%d", sw_cb1_called, sw_cb2_called); + _intg_cleanup(); +} + +static int run_stage4(void) { + printf("--- Stage 4: randomized 2-instance sync (50 iters) ---\n"); + stage_failures = 0; + test_randomized_two(); + return stage_failures; +} + +static int run_stage5(void) { + printf("--- Stage 5: randomized 3-instance star sync (30 iters) ---\n"); + stage_failures = 0; + test_randomized_three(); + return stage_failures; +} + +static int run_stage6(void) { + printf("--- Stage 6: protocol edge cases ---\n"); + stage_failures = 0; + test_cancel(); + test_start_overwrite(); + return stage_failures; +} + +int main(void) { + printf("=== test_merkle_sync ===\n"); + debug_config_init(); + debug_set_level(DEBUG_LEVEL_ERROR); + + if (run_stage1() != 0) { printf("=== Stage 1 FAILED ===\n"); return 1; } + printf("=== Stage 1 PASSED (%d tests) ===\n\n", tests_run); + + if (run_stage2() != 0) { printf("=== Stage 2 FAILED ===\n"); return 1; } + printf("=== Stage 2 PASSED (%d tests) ===\n\n", tests_run); + + if (run_stage3() != 0) { printf("=== Stage 3 FAILED ===\n"); return 1; } + printf("=== Stage 3 PASSED (%d tests) ===\n\n", tests_run); + + if (run_stage4() != 0) { printf("=== Stage 4 FAILED ===\n"); return 1; } + printf("=== Stage 4 PASSED (%d tests) ===\n\n", tests_run); + + if (run_stage5() != 0) { printf("=== Stage 5 FAILED ===\n"); return 1; } + printf("=== Stage 5 PASSED (%d tests) ===\n\n", tests_run); + + if (run_stage6() != 0) { printf("=== Stage 6 FAILED ===\n"); return 1; } + printf("=== Stage 6 PASSED (%d tests) ===\n\n", tests_run); + + printf("=== ALL PASS (%d tests) ===\n", tests_run); + return 0; +} diff --git a/tools/chatgui/transport/member_sync.c b/tools/chatgui/transport/member_sync.c index 6ea08ec1..e1341492 100644 --- a/tools/chatgui/transport/member_sync.c +++ b/tools/chatgui/transport/member_sync.c @@ -95,8 +95,8 @@ static int _member_update_bucket_hash(void* ctx, const char* ns, uint8_t level, DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: bucket_hash ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix64); char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl)); - int mask_shift = 63 - (int)level * 5; - uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; + int mask_shift = 64 - (int)level * 5; + uint64_t mask = (mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX; char sql[512]; snprintf(sql, sizeof(sql), "SELECT node_id, x25519_pubkey, ed25519_pubkey," @@ -136,8 +136,8 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level, DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: get_items ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix); char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl)); - int mask_shift = 63 - (int)level * 5; - uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; + int mask_shift = 64 - (int)level * 5; + uint64_t mask = (mask_shift >= 0 && mask_shift < 64) ? (~0ULL << mask_shift) : UINT64_MAX; char sql[512]; snprintf(sql, sizeof(sql), "SELECT node_id, x25519_pubkey, ed25519_pubkey," diff --git a/tools/chatgui/transport/merkle_sync.c b/tools/chatgui/transport/merkle_sync.c index 58e4673d..9d868307 100644 --- a/tools/chatgui/transport/merkle_sync.c +++ b/tools/chatgui/transport/merkle_sync.c @@ -13,11 +13,20 @@ #define MS_ID "merkle_sync" #define MS_BG_INTERVAL_MS 100 - -static struct merkle_sync* g_merkle = NULL; +#define MS_PENDING_MAX 128 struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; }; +enum { MS_PEND_WAITING = 0, MS_PEND_RESOLVED = 1 }; + +struct ms_pending { + uint8_t level; + uint64_t prefix; + int8_t parent; + int sub_count; + uint8_t state; +}; + struct ms_session { struct ms_session* next; char ns[64]; @@ -26,6 +35,8 @@ struct ms_session { uint8_t synced; merkle_sync_done_cb done_cb; void* cb_arg; + struct ms_pending pending[MS_PENDING_MAX]; + int pending_count; }; struct merkle_sync { @@ -49,7 +60,8 @@ static int _send_msg(struct merkle_sync* ms, uint64_t peer, const uint8_t* paylo /* ── Prefix arithmetic ── */ uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) { - int shift = 63 - (int)level * 5; if (shift < 0) shift = 0; + if (level == 0) return 0; + int shift = 64 - (int)level * 5; if (shift < 0) shift = 0; return (key >> shift) << shift; } @@ -128,8 +140,8 @@ static int _recompute_bucket(struct merkle_sync* ms, const char* ns, } void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key) { - struct merkle_sync* ms = g_merkle; - if (!ms || !ms->initialized) return; + struct merkle_sync* ms = inst->msync; + if (!ms || !ms->initialized) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: recompute_path called but msync not initialized", MS_ID); return; } DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: ns=%s key=%016llx", MS_ID, ns, (unsigned long long)key); for (uint8_t level = 1; level <= MT_MAX_LEVEL; level++) _recompute_bucket(ms, ns, level, merkle_sync_level_prefix(key, level)); @@ -165,13 +177,8 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns, if (level >= MT_MAX_LEVEL) return 0; sqlite3* db = _db(ms->inst); if (!db) return -1; - int next_shift; - if (level == 0) { - next_shift = 63 - 5; - } else { - next_shift = 63 - ((int)level + 1) * 5; if (next_shift < 0) next_shift = 0; - } + int next_shift = 64 - ((int)level + 1) * 5; if (next_shift < 0) next_shift = 0; sqlite3_stmt* stmt = NULL; int query_level = (int)(level + 1); @@ -331,6 +338,13 @@ static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, c return NULL; } +static int _find_pending(struct ms_session* s, uint8_t level, uint64_t prefix) { + for (int i = 0; i < s->pending_count; i++) + if (s->pending[i].level == level && s->pending[i].prefix == prefix) + return i; + return -1; +} + static void _session_done(struct ms_session* s, int result) { if (!s->active && s->synced) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session ALREADY DONE peer=%016llx ns=%s — SKIP second call", MS_ID, (unsigned long long)s->peer, s->ns); return; } s->active = 0; @@ -344,7 +358,7 @@ static void _session_done(struct ms_session* s, int result) { /* ── Recv handlers ── */ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, - const uint8_t* pl, size_t plen) { + const uint8_t* pl, size_t plen, int parent_idx) { if (plen < 3) return; uint8_t level = pl[0]; uint8_t pb = pl[1]; uint64_t prefix = _prefix_read(pl + 2, pb); @@ -385,21 +399,22 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns if ((local_bm & (1u << i)) && (remote_bm & (1u << i))) { int rh_idx = 0; for (int j = 0; j < i; j++) if (remote_bm & (1u << j)) rh_idx++; - if (memcmp(rh + (size_t)rh_idx * MT_HASH_SIZE, lh[i], MT_HASH_SIZE) == 0) - differs &= ~(1u << i); + if (memcmp(rh + (size_t)rh_idx * MT_HASH_SIZE, lh[i], MT_HASH_SIZE) != 0) + differs |= (1u << i); } } DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_hashes peer=%016llx ns=%s level=%d remote_bm=%08x local_bm=%08x differs=%08x is_data=%d", MS_ID, (unsigned long long)peer, ns, level, remote_bm, local_bm, differs, is_data); - if (differs == 0) { + if (differs == 0 && level == 0) { _send_hashes(ms, peer, ns, 0, 0, 0, 1); _session_done(s, MT_OK); return; } + if (differs == 0) return; - int next_shift = 63 - ((int)level + 1) * 5; + int next_shift = 64 - ((int)level + 1) * 5; struct bucket_entry requests[MT_BUCKETS]; int rcount = 0; for (int i = 0; i < MT_BUCKETS && rcount < 32; i++) { if (!(differs & (1u << i))) continue; @@ -411,6 +426,18 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns } if (rcount > 0) { + if (s) { + for (int i = 0; i < rcount && s->pending_count < MS_PENDING_MAX; i++) { + int idx = s->pending_count++; + s->pending[idx].level = requests[i].level; + s->pending[idx].prefix = requests[i].prefix; + s->pending[idx].parent = (int8_t)parent_idx; + s->pending[idx].sub_count = 0; + s->pending[idx].state = MS_PEND_WAITING; + } + if (parent_idx >= 0 && parent_idx < MS_PENDING_MAX) + s->pending[parent_idx].sub_count++; + } uint8_t ch_len = (uint8_t)strlen(ns); size_t rs = 1 + 1 + ch_len + 1; for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; @@ -448,8 +475,6 @@ static void _handle_request(struct merkle_sync* ms, uint64_t peer, const char* n uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t off = 0; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: handle_request peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count); - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_request peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count); - struct bucket_entry buckets[32]; int bc = 0; for (uint8_t i = 0; i < count && bc < 32; i++) { if (off + 2 > plen - 1) break; @@ -470,10 +495,6 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t rem = plen - 1; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: handle_batch peer=%016llx ns=%s count=%d len=%zu", MS_ID, (unsigned long long)peer, ns, count, plen - 1); - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_batch peer=%016llx ns=%s count=%d len=%zu", MS_ID, (unsigned long long)peer, ns, count, plen - 1); - - int all_terminal = 1; - for (uint8_t i = 0; i < count && rem >= 3; i++) { uint8_t lvl = bp[0]; uint8_t pb_i = bp[1]; rem -= 2; bp += 2; if (rem < pb_i + 1) break; @@ -485,14 +506,17 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, if (rem < dlen) break; ms->ops->apply_items(ms->data_ctx, ns, bp, dlen); bp += dlen; rem -= dlen; + struct ms_session* sb = _session_find(ms, peer, ns); + if (sb) { int pidx = _find_pending(sb, lvl, pr); if (pidx >= 0) sb->pending[pidx].state = MS_PEND_RESOLVED; } } else if (!is_data && rem >= 4) { - all_terminal = 0; + struct ms_session* sc = _session_find(ms, peer, ns); + int pidx = sc ? _find_pending(sc, lvl, pr) : -1; + if (pidx >= 0) sc->pending[pidx].state = MS_PEND_RESOLVED; uint8_t sub_pl[4096]; size_t sub_len = 0; - uint8_t next_lvl = (uint8_t)(lvl < MT_MAX_LEVEL ? lvl + 1 : lvl); - sub_pl[sub_len++] = next_lvl; - sub_pl[sub_len++] = merkle_sync_prefix_bytes(next_lvl); - _prefix_write(sub_pl + sub_len, pr, merkle_sync_prefix_bytes(next_lvl)); - sub_len += merkle_sync_prefix_bytes(next_lvl); + sub_pl[sub_len++] = lvl; + sub_pl[sub_len++] = merkle_sync_prefix_bytes(lvl); + _prefix_write(sub_pl + sub_len, pr, merkle_sync_prefix_bytes(lvl)); + sub_len += merkle_sync_prefix_bytes(lvl); sub_pl[sub_len++] = 0; /* is_data=0 */ uint32_t bm; memcpy(&bm, bp, 4); bp += 4; rem -= 4; memcpy(sub_pl + sub_len, &bm, 4); sub_len += 4; @@ -502,12 +526,17 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, sub_len += (size_t)nh * MT_HASH_SIZE; bp += (size_t)nh * MT_HASH_SIZE; rem -= (size_t)nh * MT_HASH_SIZE; } - _handle_hashes(ms, peer, ns, sub_pl, sub_len); + _handle_hashes(ms, peer, ns, sub_pl, sub_len, pidx); } } - if (all_terminal) { + { struct ms_session* s = _session_find(ms, peer, ns); - if (s && s->active) { + if (!s) return; + int all_done = (s->pending_count > 0); + for (int i = 0; i < s->pending_count; i++) { + if (s->pending[i].state != MS_PEND_RESOLVED) { all_done = 0; break; } + } + if (all_done && s->active) { _send_hashes(ms, peer, ns, 0, 0, 0, 1); _session_done(s, MT_OK); } @@ -519,7 +548,9 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } return; } - if (!g_merkle || !g_merkle->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } + struct UTUN_INSTANCE* inst = conn ? conn->instance : NULL; + struct merkle_sync* ms = inst ? inst->msync : NULL; + if (!ms || !ms->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } uint64_t peer = conn ? conn->peer_node_id : 0; const uint8_t* d = entry->dgram; size_t dlen = entry->len; @@ -533,16 +564,16 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { size_t plen = dlen - 3 - ch_len; switch (type) { - case 0x01: _handle_hashes(g_merkle, peer, ns, pl, plen); break; - case 0x02: _handle_request(g_merkle, peer, ns, pl, plen); break; - case 0x03: _handle_batch(g_merkle, peer, ns, pl, plen); break; + case 0x01: _handle_hashes(ms, peer, ns, pl, plen, -1); break; + case 0x02: _handle_request(ms, peer, ns, pl, plen); break; + case 0x03: _handle_batch(ms, peer, ns, pl, plen); break; case 0x04: - if (plen >= 9 && g_merkle->ops->apply_update) { + if (plen >= 9 && ms->ops->apply_update) { uint64_t key; memcpy(&key, pl, 8); uint8_t utype = pl[8]; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: recv ITEM_UPDATE key=%016llx type=%02x from=%016llx ns=%s", MS_ID, (unsigned long long)key, utype, (unsigned long long)peer, ns); - g_merkle->ops->apply_update(g_merkle->data_ctx, ns, key, utype, pl + 9, plen - 9); + ms->ops->apply_update(ms->data_ctx, ns, key, utype, pl + 9, plen - 9); } break; } @@ -554,7 +585,7 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key, uint8_t type, const uint8_t* data, size_t len) { - struct merkle_sync* ms = g_merkle; + struct merkle_sync* ms = inst ? inst->msync : NULL; if (!ms || !ms->initialized || !ns || !data) return; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: push_update ns=%s key=%016llx type=%02x len=%zu", MS_ID, ns, (unsigned long long)key, type, len); @@ -593,7 +624,7 @@ void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns, /* ── Background consistency check ── */ int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) { - struct merkle_sync* ms = g_merkle; + struct merkle_sync* ms = inst ? inst->msync : NULL; if (!ms || !ms->initialized) return -1; sqlite3* db = _db(inst); if (!db || !ns) return -1; @@ -654,7 +685,7 @@ int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, if (!ms) return -1; ms->inst = inst; ms->svc_id = svc_id; ms->ops = ops; ms->data_ctx = data_ctx; ms->initialized = 1; - g_merkle = ms; + inst->msync = ms; _ensure_table(ms); @@ -668,9 +699,10 @@ int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, } void merkle_sync_destroy(struct UTUN_INSTANCE* inst) { - struct merkle_sync* ms = g_merkle; - if (!ms || !inst) return; - ms->initialized = 0; g_merkle = NULL; + if (!inst) return; + struct merkle_sync* ms = inst->msync; + if (!ms) return; + ms->initialized = 0; inst->msync = NULL; etcp_unbind(inst, ms->svc_id); @@ -687,7 +719,7 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst) { int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns, merkle_sync_done_cb done_cb, void* arg) { - struct merkle_sync* ms = g_merkle; + struct merkle_sync* ms = inst ? inst->msync : NULL; if (!inst || !ns || !ms || !ms->initialized) return -1; _ensure_table(ms); @@ -702,6 +734,7 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, ms->sessions = s; } else { s->synced = 0; + s->pending_count = 0; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: start OVERWRITE peer=%016llx ns=%s old_cb=%p old_arg=%p → new_cb=%p new_arg=%p", MS_ID, (unsigned long long)peer, ns, (void*)s->done_cb, s->cb_arg, (void*)done_cb, arg); @@ -714,7 +747,7 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, } void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns) { - struct merkle_sync* ms = g_merkle; + struct merkle_sync* ms = inst ? inst->msync : NULL; if (!ms || !ns) return; DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: cancel peer=%016llx ns=%s", MS_ID, (unsigned long long)peer, ns); struct ms_session** p = &ms->sessions;