From ae80dae8ee2a81fe967c7650b11a8e2e6ac7b5c2 Mon Sep 17 00:00:00 2001 From: evgeny Date: Sun, 9 Aug 2026 19:39:41 +0300 Subject: [PATCH] =?UTF-8?q?tests:=20add=20Stage=207=20merkle=20stress=20te?= =?UTF-8?q?st=20=E2=80=94=2013=20instances,=205s=20spam,=20hash=20converge?= =?UTF-8?q?nce?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 10 spammers (1-30ms random) + 2 observers (1-10ms), star topology via hub. 5s of concurrent modifications, then verify all 13 merkle root hashes are identical and tree is internally consistent (recompute vs stored). --- tests/test_merkle_sync.c | 340 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 340 insertions(+) diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c index 8ba688b4..c451455c 100644 --- a/tests/test_merkle_sync.c +++ b/tests/test_merkle_sync.c @@ -974,6 +974,336 @@ static int run_stage5(void) { return stage_failures; } +/* ── Stage 7: stress test — 10 spammers + 2 A↔B observers, star topology ── */ + +#define STRESS_N 13 /* 1 hub + 10 spammers + 2 observers */ +#define STRESS_HUB 0 +#define STRESS_SPAM_BEGIN 1 +#define STRESS_SPAM_END 10 /* indices 1..10 */ +#define STRESS_OBS_A 11 +#define STRESS_OBS_B 12 +#define STRESS_SPAM_MS 5000 +#define STRESS_SYNC_TB 200000 /* 20s timeout for final sync */ +#define STRESS_LINK_TB 300000 /* 30s for all links up */ + +struct str_state { + struct UASYNC* ua; + struct UTUN_INSTANCE* inst[STRESS_N]; + struct ms_data data[STRESS_N]; + struct ms_test_ctx ctx[STRESS_N]; + void* global_timeout; + int phase; /* 0=running, 1=stopping, 2=timeout */ + int spam_active; + int sync_pending; /* atomic: count of outstanding sync ops */ +}; + +static struct str_state* gs = NULL; + +static void _str_cleanup(void); +static int _str_cond_sync_done(void); +static int _str_cond_links_up(void); +static void _str_wait_sync_quiesce(void); + +static void _str_sync_done(uint64_t peer, const char* ns, int result, void* arg) { + (void)peer; (void)ns; (void)result; + struct str_state* s = (struct str_state*)arg; + if (s) __sync_fetch_and_sub(&s->sync_pending, 1); +} + +static void _str_spam_cb(void* arg); + +static void _str_spam_schedule(struct str_state* s, int idx) { + if (!s->spam_active) return; + int delay_tb = ((rand() % 30) + 1) * 10; /* 1-30ms → 10-300 timebase units */ + uasync_set_timeout(s->ua, (uint32_t)delay_tb, (void*)(intptr_t)idx, _str_spam_cb, "spam"); +} + +static void _str_spam_cb(void* arg) { + if (!gs || !gs->spam_active) return; + int idx = (int)(intptr_t)arg; + if (idx < STRESS_SPAM_BEGIN || idx > STRESS_SPAM_END) { _str_spam_schedule(gs, idx); return; } + + /* spammer: bump own member, sync to hub */ + uint64_t key = ((uint64_t)idx) << 60; + _data_insert(&gs->data[idx], key, (uint32_t)(gs->data[idx].count + 1)); + merkle_sync_recompute_path(gs->inst[idx], "stress", key); + __sync_fetch_and_add(&gs->sync_pending, 1); + merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", _str_sync_done, gs); + + _str_spam_schedule(gs, idx); +} + +static void _str_obs_spam_cb(void* arg) { + if (!gs || !gs->spam_active) return; + int idx = (int)(intptr_t)arg; + + __sync_fetch_and_add(&gs->sync_pending, 1); + merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", _str_sync_done, gs); + + int delay_tb = ((rand() % 10) + 1) * 10; /* 1-10ms */ + uasync_set_timeout(gs->ua, (uint32_t)delay_tb, arg, _str_obs_spam_cb, "obs_spam"); +} + +static void _str_global_timeout(void* arg) { + (void)arg; + if (gs) gs->phase = 2; +} + +static int _str_cond_links_up(void) { + if (!gs) return 0; + int total_links = 0; + for (int i = 0; i < STRESS_N; i++) { + if (!gs->inst[i] || !gs->inst[i]->connections) return 0; + struct ll_entry* e = gs->inst[i]->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) total_links++; l = l->next; } + e = e->next; + } + } + /* hub should have 12 links (one per spoke), spokes have 1 each */ + return total_links >= STRESS_N * 2 - 2; /* 12 hub links + 12 spoke links = 24 total */ +} + +static int _str_wait_for(const char* desc, int (*cond)(void), int timeout_tb) { + struct str_state* s = gs; + uint64_t start = get_time_tb(); + while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && s->phase == 0) + uasync_poll(s->ua, 1); + if (cond()) return 1; + printf(" STRESS TIMEOUT: %s\n", desc); s->phase = 2; + return 0; +} + +static int _str_init(void) { + struct str_state* s = u_calloc(1, sizeof(*s)); + if (!s) return -1; + gs = s; + s->spam_active = 0; + s->phase = 0; + s->sync_pending = 0; + + /* --- temp dir --- */ + char tmpld[128] = "/tmp/utun_str_XXXXXX"; + if (test_mkdtemp(tmpld) != 0) { printf(" mkdtemp failed\n"); u_free(s); gs = NULL; return -1; } + + /* --- create configs --- */ + int base_port = 44000 + (getpid() % 10000); + char cfgs[STRESS_N][256]; int ports[STRESS_N]; + for (int i = 0; i < STRESS_N; i++) ports[i] = base_port + i; + + utun_instance_set_tun_init_enabled(0); + s->ua = uasync_create(); + if (!s->ua) goto fail; + + /* hub config */ + snprintf(cfgs[0], sizeof(cfgs[0]), "%s/cfg_0.conf", tmpld); + { + char dbp[256]; snprintf(dbp, sizeof(dbp), "%s/db_0", tmpld); utun_mkdir(dbp, 0755); + _write_cfg(cfgs[0], + "[global]\nmy_node_id=0xAAAAAAAAAAAAAAAA\ntun_ip=10.240.0.1/24\ntun_ifname=tun240\nkeepalive_adaptive=0\ndb_path=%s\n\n" + "[server:s0]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[allowed_keys]\nallow_all=1\n", dbp, ports[0]); + config_ensure_keys_and_node_id(cfgs[0]); + } + + char* hub_pub = _get_pubkey(cfgs[0]); if (!hub_pub) goto fail; + /* spoke configs */ + for (int i = 1; i < STRESS_N; i++) { + snprintf(cfgs[i], sizeof(cfgs[i]), "%s/cfg_%d.conf", tmpld, i); + char dbp[256]; snprintf(dbp, sizeof(dbp), "%s/db_%d", tmpld, i); utun_mkdir(dbp, 0755); + char nid[32]; snprintf(nid, sizeof(nid), "BBBBBBBBBBBB%02d00", i); + _write_cfg(cfgs[i], + "[global]\nmy_node_id=0x%s\ntun_ip=10.240.0.%d/24\ntun_ifname=tun240\nkeepalive_adaptive=0\ndb_path=%s\n\n" + "[server:s%d]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[client:to_hub]\nkeepalive=1\npeer_public_key=%s\nlink=s%d:127.0.0.1:%d\n\n" + "[allowed_keys]\nallow_all=1\n", + nid, i + 1, dbp, i, ports[i], hub_pub, i, ports[0]); + config_ensure_keys_and_node_id(cfgs[i]); + } + u_free(hub_pub); + + /* --- create instances --- */ + for (int i = 0; i < STRESS_N; i++) { + s->inst[i] = utun_instance_create(s->ua, cfgs[i]); + if (!s->inst[i] || utun_instance_init(s->inst[i]) != 0) { + printf(" FAIL: instance %d init\n", i); + goto fail; + } + } + printf(" %d instances created, waiting for links...\n", STRESS_N); + + s->global_timeout = uasync_set_timeout(s->ua, STRESS_LINK_TB, NULL, _str_global_timeout, "str_timeout"); + if (!_str_wait_for("links up", _str_cond_links_up, STRESS_LINK_TB)) { printf(" links never came up\n"); goto fail; } + printf(" all links UP\n"); + uasync_cancel_timeout(s->ua, s->global_timeout); s->global_timeout = NULL; + + /* --- init merkle_sync on all --- */ + for (int i = 0; i < STRESS_N; i++) { + memset(&s->data[i], 0, sizeof(s->data[i])); + s->ctx[i].data = &s->data[i]; + s->ctx[i].inst = s->inst[i]; + s->ctx[i].tag = (char)('A' + i); + if (merkle_sync_init(s->inst[i], MS_SVC_ID, &g_test_ops, &s->ctx[i]) != 0) { + printf(" FAIL: merkle_sync_init %d\n", i); goto fail; + } + } + printf(" merkle_sync inited on all\n"); + + /* --- seed each instance with its own member; sync to hub --- */ + for (int i = 1; i < STRESS_N; i++) { + uint64_t key = ((uint64_t)i) << 60; + _data_insert(&s->data[i], key, (uint32_t)i); + merkle_sync_recompute_path(s->inst[i], "stress", key); + } + /* initial sync: each spoke → hub */ + s->sync_pending = STRESS_N - 1; + for (int i = 1; i < STRESS_N; i++) + merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, "stress", _str_sync_done, s); + _str_wait_for("initial sync", _str_cond_sync_done, STRESS_SYNC_TB); + + /* observers: sync with hub too */ + + return 0; + +fail: + _str_cleanup(); + return -1; +} + +static int _str_cond_sync_done(void) { + return gs ? __sync_fetch_and_add(&gs->sync_pending, 0) <= 0 : 0; +} + +static int _str_cond_sync_done_strict(void) { + return gs ? __sync_fetch_and_add(&gs->sync_pending, 0) <= 0 : 0; +} + +static void _str_wait_sync_quiesce(void) { + struct str_state* s = gs; + uint64_t start = get_time_tb(); + while (s->sync_pending > 0 && (get_time_tb() - start) < (uint64_t)STRESS_SYNC_TB) + uasync_poll(s->ua, 1); +} + +static void _str_cleanup(void) { + struct str_state* s = gs; + if (!s) return; + s->spam_active = 0; + if (s->global_timeout && s->ua) { uasync_cancel_timeout(s->ua, s->global_timeout); s->global_timeout = NULL; } + for (int i = STRESS_N - 1; i >= 0; i--) { + if (!s->inst[i]) continue; + merkle_sync_destroy(s->inst[i]); + s->inst[i]->running = 0; + utun_instance_destroy(s->inst[i]); + s->inst[i] = NULL; + } + if (s->ua) { uasync_destroy(s->ua, 0); s->ua = NULL; } + u_free(s); + gs = NULL; +} + +static void test_stress_spam(void) { + TEST("stress: 10 spammers + 2 obs, 5s spam, verify merkle roots"); + if (_str_init() != 0) { FAIL("init failed"); _str_cleanup(); return; } + + struct str_state* s = gs; + + /* --- start spam --- */ + s->spam_active = 1; + srand(12345); + for (int i = STRESS_SPAM_BEGIN; i <= STRESS_SPAM_END; i++) { + /* stagger initial delays */ + int d = (rand() % 30) + 1; + uasync_set_timeout(s->ua, (uint32_t)(d * 10), (void*)(intptr_t)i, _str_spam_cb, "spam"); + } + /* observers sync with hub aggressively */ + uasync_set_timeout(s->ua, 10, (void*)(intptr_t)STRESS_OBS_A, _str_obs_spam_cb, "obsA"); + uasync_set_timeout(s->ua, 15, (void*)(intptr_t)STRESS_OBS_B, _str_obs_spam_cb, "obsB"); + + /* --- run for 5 seconds --- */ + printf(" spamming for %dms...\n", STRESS_SPAM_MS); + uint64_t start = get_time_tb(); + while ((get_time_tb() - start) < (uint64_t)(STRESS_SPAM_MS * 10) && s->phase == 0) + uasync_poll(s->ua, 1); + + /* --- stop all spam --- */ + s->spam_active = 0; + + /* let in-flight syncs settle: poll for up to 10s */ + printf(" sync_pending=%d, waiting for convergence...\n", s->sync_pending); + { + uint64_t settle_start = get_time_tb(); + int last_pending = s->sync_pending; + while ((get_time_tb() - settle_start) < (uint64_t)(STRESS_SYNC_TB)) { + uasync_poll(s->ua, 1); + int cur = __sync_fetch_and_add(&s->sync_pending, 0); + if (cur == 0) break; + if (cur != last_pending) { last_pending = cur; settle_start = get_time_tb(); } + } + } + printf(" converged: sync_pending=%d phase=%d\n", s->sync_pending, s->phase); + + if (s->phase == 2) { FAIL("global timeout"); _str_cleanup(); return; } + + /* --- final explicit sync from each spoke → hub, then wait --- */ + for (int i = 1; i < STRESS_N; i++) + merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); + /* poll for data propagation */ + for (int i = 0; i < 500; i++) uasync_poll(s->ua, 5); + + /* --- verify: tree sizes match pairwise --- */ + int ok = 1; + int na = _db_count_rows(s->inst[0]->topo_sqlite_db, "stress"); + for (int i = 1; i < STRESS_N && ok; i++) { + int nb = _db_count_rows(s->inst[i]->topo_sqlite_db, "stress"); + if (na != nb) { printf(" tree size mismatch: hub=%d inst[%d]=%d\n", na, i, nb); ok = 0; } + } + + /* --- verify: all root (level=1, prefix=0) hashes identical --- */ + const uint8_t* root0 = merkle_sync_get_hash(s->inst[0], "stress", 1, 0); + uint8_t root_copy[MT_HASH_SIZE]; memcpy(root_copy, root0, MT_HASH_SIZE); + for (int i = 1; i < STRESS_N && ok; i++) { + const uint8_t* ri = merkle_sync_get_hash(s->inst[i], "stress", 1, merkle_sync_level_prefix(0, 1)); + if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { printf(" ROOT HASH MISMATCH inst[%d]\n", i); ok = 0; } + } + + if (!ok) { FAIL("hash mismatch"); _str_cleanup(); return; } + + /* --- verify: recompute consistency for each instance --- */ + for (int i = 0; i < STRESS_N && ok; i++) { + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(s->inst[i]->topo_sqlite_db, + "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", + -1, &st, NULL); + sqlite3_bind_text(st, 1, "stress", -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(&s->ctx[i], "stress", 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(s->inst[i], "stress", lv, pf); + uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE); + if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { + printf(" consistency fail inst[%d] L%d P%016llx\n", i, lv, (unsigned long long)pf); + ok = 0; + } + } + sqlite3_finalize(st); + } + if (!ok) { FAIL("tree consistency"); _str_cleanup(); return; } + + /* print summary */ + printf(" tree rows: %d data items per node:", na); + for (int i = 0; i < STRESS_N && i < 6; i++) printf(" %d", s->data[i].count); + printf("..\n"); + PASS(); + _str_cleanup(); +} static int run_stage6(void) { printf("--- Stage 6: protocol edge cases ---\n"); stage_failures = 0; @@ -982,6 +1312,13 @@ static int run_stage6(void) { return stage_failures; } +static int run_stage7(void) { + printf("--- Stage 7: stress test (10 spammers + 2 A↔B obs, 5s) ---\n"); + stage_failures = 0; + test_stress_spam(); + return stage_failures; +} + int main(void) { printf("=== test_merkle_sync ===\n"); debug_config_init(); @@ -1005,6 +1342,9 @@ int main(void) { if (run_stage6() != 0) { printf("=== Stage 6 FAILED ===\n"); return 1; } printf("=== Stage 6 PASSED (%d tests) ===\n\n", tests_run); + if (run_stage7() != 0) { printf("=== Stage 7 FAILED ===\n"); return 1; } + printf("=== Stage 7 PASSED (%d tests) ===\n\n", tests_run); + printf("=== ALL PASS (%d tests) ===\n", tests_run); return 0; }