#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, uint64_t from_peer, 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; uint8_t relay_buf[65536]; uint16_t relay_count = 0; size_t relay_off = 2; /* [count:2][items] */ 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; /* check if item is new or different */ int found = 0; for (int j = 0; j < d->count; j++) { if (d->items[j].key == key) { if (d->items[j].val != val) { d->items[j].val = val; changed = 1; } found = 1; break; } } if (!found) { _data_insert(d, key, val); changed = 1; } if (changed && relay_off + 12 < sizeof(relay_buf)) { memcpy(relay_buf + relay_off, &key, 8); relay_off += 8; memcpy(relay_buf + relay_off, &val, 4); relay_off += 4; relay_count++; } } 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); } /* relay changed items to other SYNCED sessions */ if (relay_count > 0) { memcpy(relay_buf, &relay_count, 2); merkle_sync_broadcast(tc->inst, ns, from_peer, relay_buf, relay_off); } } 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 PHASE3_TIMEOUT_TB 200000 #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) uasync_poll(i_ua, POLL_MS); if (cond()) return 1; 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 A=%016llx B=%016llx\n", (unsigned long long)i_a->node_id, (unsigned long long)i_b->node_id); 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 — 25 iterations"); srand(42); int iter; for (iter = 0; iter < 25 && 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*5, NULL, _intg_timeout, "ms_intg_timeout"); if (!_wait_for("links up", _cond_links_up, PHASE3_TIMEOUT_TB)) { _intg_cleanup(); return -1; } /* let connections settle */ for (int i = 0; i < 20; i++) uasync_poll(i_ua, 5); _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 A=%016llx B=%016llx C=%016llx\n", (unsigned long long)i_a->node_id, (unsigned long long)i_b->node_id, (unsigned long long)i_c->node_id); 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; } static void test_randomized_three(void) { TEST("randomized 3-instance star — 5 iterations"); srand(1234); int iter; for (iter = 0; iter < 5 && 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; i_phase = 0; /* B→A: B syncs with A, gets A's data, B: SYNCED */ merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); if (!_wait_for("sync B", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout B iter %d", iter); _intg_cleanup(); break; } done_sync = 0; i_phase = 0; /* C→A: C syncs with A, gets A∪B, A broadcasts C's new items to B (relay) */ merkle_sync_start(i_c, i_a->node_id, "rnd", _on_sync_done, NULL); if (!_wait_for("sync C", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; } /* wait for broadcast relay */ for (int k = 0; k < 100; k++) uasync_poll(i_ua, 10); 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 == 5) 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; }