diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c index c451455c..08f68fbe 100644 --- a/tests/test_merkle_sync.c +++ b/tests/test_merkle_sync.c @@ -856,7 +856,7 @@ static void test_randomized_two(void) { if (!ok) { FAIL("tree consistency iter %d", iter); _intg_cleanup(); break; } _intg_cleanup(); } - if (iter == 50) PASS(); + if (iter == 25) PASS(); } /* ── Stage 5: randomized 3-instance star-topology sync ── */ @@ -985,16 +985,26 @@ static int run_stage5(void) { #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 */ +#define STRESS_MAX_ROUNDS 5 +#define STRESS_DRAIN_TB 2000 /* 200ms drain after spam stop */ +#define STRESS_ROUND_GAP_TB 500 /* 50ms gap between final rounds */ +#define STRESS_SAFETY_TB 400000 /* 40s safety for spam+converge */ 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 */ + void* global_timeout; /* setup: links-up timeout */ + void* safety_timer; /* spam+converge safety */ + void* stop_timer; /* spam stop */ + void* drain_timer; /* drain after spam stop */ + void* round_timer; /* gap between final rounds */ + int phase; /* 0=running, 1=pass, 2=fail */ int spam_active; - int sync_pending; /* atomic: count of outstanding sync ops */ + int sync_pending; /* initial sync only (balanced) */ + int final_pending; /* outstanding final-round syncs */ + int final_round; /* final round number */ }; static struct str_state* gs = NULL; @@ -1002,7 +1012,14 @@ 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_stop_spam(void* arg); +static void _str_drain_cb(void* arg); +static void _str_round_cb(void* arg); +static void _str_start_final_round(void); +static void _str_final_done(uint64_t peer, const char* ns, int result, void* arg); +static void _str_check_roots(void); +static void _str_verify(void); +static void _str_fail_timeout(void* arg); static void _str_sync_done(uint64_t peer, const char* ns, int result, void* arg) { (void)peer; (void)ns; (void)result; @@ -1027,8 +1044,7 @@ static void _str_spam_cb(void* arg) { 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); + merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); _str_spam_schedule(gs, idx); } @@ -1037,8 +1053,7 @@ 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); + merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); 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"); @@ -1176,22 +1191,15 @@ 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; } + if (s->safety_timer && s->ua) { uasync_cancel_timeout(s->ua, s->safety_timer); s->safety_timer = NULL; } + if (s->stop_timer && s->ua) { uasync_cancel_timeout(s->ua, s->stop_timer); s->stop_timer = NULL; } + if (s->drain_timer && s->ua) { uasync_cancel_timeout(s->ua, s->drain_timer); s->drain_timer = NULL; } + if (s->round_timer && s->ua) { uasync_cancel_timeout(s->ua, s->round_timer); s->round_timer = NULL; } for (int i = STRESS_N - 1; i >= 0; i--) { if (!s->inst[i]) continue; merkle_sync_destroy(s->inst[i]); @@ -1204,74 +1212,60 @@ static void _str_cleanup(void) { 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); +static void _str_stop_spam(void* arg) { + (void)arg; + if (!gs) return; + gs->stop_timer = NULL; + gs->spam_active = 0; + printf(" spam stopped\n"); + /* короткий дренаж — даём in-flight спам-синкам утихнуть перед финальными раундами */ + gs->drain_timer = uasync_set_timeout(gs->ua, STRESS_DRAIN_TB, NULL, _str_drain_cb, "str_drain"); +} - if (s->phase == 2) { FAIL("global timeout"); _str_cleanup(); return; } +static void _str_drain_cb(void* arg) { + (void)arg; + if (!gs) return; + gs->drain_timer = NULL; + _str_start_final_round(); +} - /* --- final explicit sync from each spoke → hub, then wait --- */ +static void _str_start_final_round(void) { + struct str_state* s = gs; + if (!s) return; + s->final_round++; + s->final_pending = STRESS_N - 1; + printf(" final round %d: syncing %d spokes -> hub\n", s->final_round, s->final_pending); 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); + merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, "stress", _str_final_done, NULL); +} - /* --- verify: tree sizes match pairwise --- */ +static void _str_final_done(uint64_t peer, const char* ns, int result, void* arg) { + (void)peer; (void)ns; (void)result; (void)arg; + if (!gs) return; + gs->final_pending--; + if (gs->final_pending <= 0) _str_check_roots(); +} + +static void _str_verify(void) { + struct str_state* s = gs; int ok = 1; - int na = _db_count_rows(s->inst[0]->topo_sqlite_db, "stress"); + + /* tree sizes match pairwise */ + int na = _db_count_rows(s->inst[STRESS_HUB]->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); + /* root (level=1, prefix=0) hashes identical */ + const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], "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)); + const uint8_t* ri = merkle_sync_get_hash(s->inst[i], "stress", 1, 0); 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 --- */ + /* 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, @@ -1295,12 +1289,82 @@ static void test_stress_spam(void) { } 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"); + + s->phase = ok ? 1 : 2; +} + +static void _str_check_roots(void) { + struct str_state* s = gs; + if (!s) return; + const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], "stress", 1, 0); + uint8_t root_copy[MT_HASH_SIZE]; + memcpy(root_copy, root0, MT_HASH_SIZE); + int converged = 1; + for (int i = 1; i < STRESS_N; i++) { + const uint8_t* ri = merkle_sync_get_hash(s->inst[i], "stress", 1, 0); + if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { converged = 0; break; } + } + if (converged) { + printf(" roots converged after %d round(s)\n", s->final_round); + _str_verify(); + } else if (s->final_round < STRESS_MAX_ROUNDS) { + s->round_timer = uasync_set_timeout(s->ua, STRESS_ROUND_GAP_TB, NULL, _str_round_cb, "str_round"); + } else { + printf(" no convergence after %d rounds\n", s->final_round); + s->phase = 2; + } +} + +static void _str_round_cb(void* arg) { + (void)arg; + if (!gs) return; + gs->round_timer = NULL; + _str_start_final_round(); +} + +static void _str_fail_timeout(void* arg) { + (void)arg; + if (!gs) return; + printf(" SAFETY TIMEOUT\n"); + gs->phase = 2; +} + +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"); + + /* --- spam stop + safety timeout: событиями, без холостых ожиданий --- */ + s->stop_timer = uasync_set_timeout(s->ua, (uint32_t)(STRESS_SPAM_MS * 10), NULL, _str_stop_spam, "str_stop"); + s->safety_timer = uasync_set_timeout(s->ua, STRESS_SAFETY_TB, NULL, _str_fail_timeout, "str_safety"); + + printf(" spamming for %dms...\n", STRESS_SPAM_MS); + /* Событийный цикл: фазы переключаются таймерами и done-коллбэками. + * poll с ограниченным таймаутом (10ms), т.к. uasync_poll(ua,-1) при уже + * просроченном ближайшем таймере уходит в бесконечный epoll_wait и не + * обрабатывает expired-таймеры. */ + while (s->phase == 0) + uasync_poll(s->ua, 100); + + if (s->phase != 1) { FAIL("convergence failed (phase=%d)", s->phase); _str_cleanup(); return; } PASS(); _str_cleanup(); }